6 Commits
8 ... 12

Author SHA1 Message Date
subochev f70a9fdc9c feat: Prometheus-метрики GET /metrics + ленивый ISO-8601-парсер для backoff
Build LLM Proxy / Build and push (release) Successful in 40s
- GET /metrics (Prometheus text 0.0.4): llm_proxy_requests_total{model,upstream,provider,result}
  (ok/4xx/429/402/5xx/net_err/cancelled), llm_proxy_upstream_inflight,
  llm_proxy_upstream_fail_streak, llm_proxy_upstream_cooling_seconds
- учёт исходов в handleChat (каждая фейловер-попытка — отдельно)
- parseIsoDuration: PnM≈n×30d, PnY≈n×365d (kotlin.time Duration.parse не принимает Y/M)
- remainingSeconds округляет остаток вверх (открытое окно ≥1s)
- доки (CONFIG.md/README/TESTING) + тесты (MetricsTest, BackoffTest)
2026-09-13 20:59:28 +03:00
subochev 0ba9b7d739 feat: экспоненциальный backoff (providers[].backoff / upstreams[].backoff)
Повторяющиеся сбои апстрима (5xx/429/402 и сетевые ошибки) увязываются
экспоненциальным откатом: первая ошибка — отдых 1s, далее 2s, 4s, … до
потолка из конфига (ISO-8601: P1M, PT15M…); успешный запрос сбрасывает
счётчик. Приоритет: upstreams[].backoff (свой счётчик на модель) над
providers[].backoff (общий счётчик на все модели провайдера).

- Backoff.kt: BackoffGuard (Mutex/@Volatile, kotlinx-datetime) +
  BackoffRegistry (скоупы провайдер/модель);
- pickFreeUpstream пропускает апстримы в откате (лог «~Ns — пропускаю»);
- все апстримы модели в откате → 503 «all upstreams cooling, retry in Ns»
  + заголовок Retry-After;
- некорректный ISO-8601 в backoff — warn + игнор поля;
- BackoffTest (8 тестов) + обновлённые тесты pickFreeUpstream;
- CONFIG.md/README/TESTING.md. Проверено: jvmTest, compileKotlinLinuxX64,
  linuxX64Test.
2026-09-13 20:33:36 +03:00
subochev adaffba885 feat: 402 (insufficient balance / лимит токенов) — фейловер на следующий апстрим
402 от провайдера — ошибка аккаунта (баланс/лимит токенов), а не запроса:
следующий апстрим в списке модели может быть платёжеспособен. Теперь при
5xx/429/402 освобождаем слот и пробуем следующего свободного; остальные 4xx
продолжаем отдавать клиенту как есть. Обновлён CONFIG.md (маршрутизация,
лог-секции). Ранее этот правка терялась (3155065 не сохранился) — восстанавливаю.
2026-09-13 20:33:32 +03:00
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
subochev 558b2f925b feat: выравнивание контракта конца SSE — дописывать data: [DONE] при штатном закрытии
Build LLM Proxy / Build and push (release) Successful in 38s
Проблема (Vikunja #71): opencode показывал «модель ещё думает» ~45 с после конца
генерации на апстриме MiniMax (модель codding-big), который НЕ присылает
терминатор data: [DONE] и просто закрывает сокет. Bifrost для custom-провайдеров
считает маркер обязательным, поэтому завершал стрим клиенту только на закрытии
keep-alive-сокета llm-proxy (Ktor CIO connectionIdleTimeoutSeconds = 45).

Правка: стрим клиенту завершается маркером data: [DONE] всегда, когда апстрим
закрыл поток штатно (в потоке был непустой finish_reason), но сам маркера не
прислал. Обрыв БЕЗ finish_reason маркер НЕ дописывает — иначе усечённый стрим
выглядит как успешный (детект Bifrost SSEStreamEndedOnMarker).

- Main.kt: SSE_DONE_MARKER, StreamEndDetector (скользящее окно 128 байт — маркер
  может разрезаться границей чтения), hasNonNullFinishReasonText (экранированные
  вхождения не считаются), hasFinishReason, streamRawWithDoneContract для ветки
  think_tags: off; в streamSseWithThinkTags — учёт sawDone/sawFinishReason и
  дописывание маркера после сброса хвостов сплиттеров.
- StreamDoneContractTest.kt: 13 тестов (детектор, разрез маркера границей чтения,
  byte-exact passthrough, отсутствие дублирования, обрыв без finish_reason).
- TESTING.md: строка про новый тест-файл + 3 пункта мутационной приёмки;
- docs/sse-done-contract.md: ТЗ и разбор замеров.

Проверено: ./gradlew clean jvmTest fatJar — 79 тестов (было 66), 0 падений;
мутационная приёмка — ослабление охраны до `if (!sawDone)` роняет
thinkStreamTruncatedNeedsNoMarker.
2026-09-11 18:26:52 +03:00
subochev 334015a2bd test: покрытие think_tags — non-stream, SSE-чанки, обвязка стрима
Build LLM Proxy / Build and push (release) Successful in 44s
- transformThinkMessage/transformThinkChunk/streamSseWithThinkTags сделаны internal
  ради тестируемости (логика не менялась)
- +24 теста: ThinkTagTransformTest (9), ThinkTagChunkTest (8), ThinkTagStreamTest (7)
- ConfigLogicTest: ассерт приоритета источников был неразличающим (провайдер и
  апстрим давали одинаковый результат) — заменён на различающиеся значения
- kotlinx-coroutines-test для тестов каналов ktor
- Gitea Actions: шаг Run tests (jvmTest) — раньше CI тесты не гонял вовсе
- TESTING.md: что покрыто + мутационная приёмка
2026-09-11 17:40:06 +03:00
19 changed files with 2260 additions and 67 deletions
+4
View File
@@ -9,6 +9,10 @@ jobs:
steps:
- name: 'Checkout'
uses: https://github.com/actions/checkout@v4
- name: 'Run tests'
uses: https://git.binom.pw/subochev/devops/build-gradle@main
with:
target: jvmTest
- name: 'Build jar'
uses: https://git.binom.pw/subochev/devops/build-gradle@main
with:
+142 -6
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,8 @@ providers:
patch: # уровень провайдера: ко всем его запросам
provider:
allow_fallbacks: false
backoff: P1M # опционально; потолок экспоненциального backoff
# на весь провайдер (ISO-8601: P1M, PT15M, P1D…)
- id: local-llama
url: "http://10.0.0.5:8080/v1"
@@ -93,6 +101,8 @@ upstreams:
patch: # уровень апстрима
provider:
ignore: [deepseek]
backoff: PT15M # опционально; собственный backoff этой модели
# (приоритет над backoff провайдера)
- id: local-qwen
provider: local-llama
@@ -218,6 +228,120 @@ 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` в истории могут и не передавать.
Поля (только на уровне **провайдера**, это свойство шлюза, а не модели):
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) |
| `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе |
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в каждом
assistant-сообщении с непустым `tool_calls` — берёт текст из `reasoning`
(если это непустая строка), иначе склеивает `reasoning_details[*].text`
(только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и
записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее
непустое поле не перезаписывается. Ничего при этом не убирается и не
переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте,
поле просто дополняется. Тронуто только assistant-сообщение с непустым
`tool_calls`; сообщения без `tool_calls` и `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`.
### 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`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -296,9 +420,13 @@ 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)
)
@Serializable
@@ -308,6 +436,8 @@ data class UpstreamConf(
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 на модель (приоритет над провайдерским)
)
@Serializable
@@ -413,7 +543,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 +596,14 @@ 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=…`).
- **Все в откате** — `503 all upstreams cooling, retry in Ns` для
`model=<витрина>` с заголовком `Retry-After: N`.
Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать
существующий парсинг логов (если он есть).
+12 -1
View File
@@ -20,10 +20,21 @@ 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`.
Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по
`model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) —
подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список
метрик — в CONFIG.md, раздел «Prometheus-метрики».
| Переменная | Default | Описание |
|---|---|---|
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
+90
View File
@@ -0,0 +1,90 @@
# Тестирование llm-proxy
Все тесты — обычные модульные, живут в `src/commonTest/kotlin/pw/binom/llmproxy/`.
## Запуск
- `./gradlew jvmTest` — прогнать все тесты;
- `./gradlew clean jvmTest fatJar` — полная сборка с нуля;
- результат смотреть в `build/test-results/jvmTest/*.xml` (атрибуты `tests`/`failures`/`errors`), потому что строки вида «N tests completed» печатаются только при падениях.
## Что покрыто
| Файл | Что проверяет |
| --- | --- |
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
| `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, неприкосновенность сообщений без `tool_calls` (в т.ч. `[]`) и `user`/`tool`, пустая строка при `emptyOk`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
| `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`. |
## Проверка качества тестов (мутационная приёмка)
Приём по шагам:
1. Забэкапить файл.
2. Внести РОВНО одну поломку в боевой код.
3. Прогнать `./gradlew cleanJvmTest jvmTest`.
4. Посмотреть XML — тест, который не упал, считается пустым.
5. Откатить (`git checkout -- <файл>`).
Обязательно: `cleanJvmTest` обязателен, иначе прогон не перезапустится.
Проверенные мутации, каждая из которых ДОЛЖНА ронять тесты:
- `transformThinkMessage` возвращает null → падают тесты non-stream;
- `transformThinkChunk` возвращает null → падают тесты чанков и стрима;
- блок финиш-чанка в `streamSseWithThinkTags` не выполняется → падает тест про удержанный хвост;
- `holdableSuffix` всегда 0 (хвост тега не удерживается) → падают тесты автомата, чанков и стрима;
- `strip` начинает отдавать рассуждения → падает тест автомата;
- `missingDoneMarker` всегда `false` → падают тесты сырого пути (`rawStreamIsByteExactAndAppendsDone`) и тест разрезанного маркера;
- `if (sawFinishReason && !sawDone)` → `if (!sawDone)` в `streamSseWithThinkTags` → падает `thinkStreamTruncatedNeedsNoMarker`;
- `if (sawFinishReason && !sawDone)` → `if (false)` в `streamSseWithThinkTags` → падает `thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker`.
Правило: боевой код нельзя подгонять под тест; если тест не проходит, неверен тест.
## Поле рассуждений для 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`.
## Известное ограничение
Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
Правка поля рассуждений — исключение: она проверена живым инстансом против настоящего Console Go (таблица выше) до релиза.
+1
View File
@@ -31,6 +31,7 @@ kotlin {
val commonTest by getting {
dependencies {
implementation(libs.kotlin.test)
implementation(libs.kotlinx.coroutines.test)
}
}
val jvmMain by getting {
+280
View File
@@ -0,0 +1,280 @@
# ТЗ: выравнивание контракта конца SSE (дописать `data: [DONE]`)
Задача-источник: Vikunja **#71** («в opencode модель додумала, идёт ожидание ~45 с»).
## 1. Что происходит сейчас (проверено живыми замерами 11.09.2026)
Цепочка: opencode → `llm.binom.pw` (Bifrost v1.5.0) → `llm-proxy:8100` → `api.minimax.io`.
- **MiniMax не присылает терминатор `data: [DONE]`.** После чанка с `finish_reason: "stop"` он просто закрывает
сокет (сразу). Проверено и с `reasoning_split`, и с `stream_options.include_usage`.
- **llm-proxy отдаёт поток корректно**, но клиент ждёт маркер: при `think_tags: off` (это как раз MiniMax-апстрим
`minimax-m3` модели `codding-big`) работает сырой байтовый passthrough — что прислал апстрим, то и ушло.
Маркера нет → нет маркера.
- **Bifrost v1.5.0** для custom-провайдеров считает `[DONE]` обязательным (`ProviderSendsDoneMarker` = true;
флага `custom_provider_config.does_not_send_done_marker` в этой версии нет). Маркера нет → его цикл чтения
не завершается, и стрим закрывается только когда умирает keep-alive-сокет llm-proxy — через **45.00 с**
(Ktor CIO `connectionIdleTimeoutSeconds = 45`). После этого Bifrost дописывает клиенту `[DONE]` и закрывает поток.
Отсюда «модель додумала, а opencode ещё думает».
- Контроль: тот же Bifrost и та же тула, апстрим со своим `[DONE]` (RouterAI `local/deepseek/deepseek-v4-flash-0731`)
— закрытие за **0.00 с**.
## 2. Требуемое поведение (контракт)
Стрим, который отдаёт llm-proxy клиенту, **всегда** заканчивается маркером `data: [DONE]\n\n`, если апстрим
завершил поток штатно (в потоке был непустой `finish_reason`), — даже когда сам апстрим маркер не прислал.
Правило ровно в двух частях:
1. Апстрим закрыл поток, и в потоке был непустой `finish_reason`, но `data: [DONE]` мы не видели →
дописать клиенту `data: [DONE]\n\n` и `flush()`.
2. Апстрим оборвался, **не** прислав `finish_reason` (усечённый/битый поток) → маркер **НЕ** дописывать,
отдать как есть.
Пункт 2 — не забыть и не «улучшить»: именно по отсутствию маркера Bifrost отличает усечённый стрим от
успешного (его детект `SSEStreamEndedOnMarker`). Если дописать маркер при обрыве без `finish_reason`,
мы превратим обрыв в «успех» и сломаем диагностику.
Обязательно: **байты уже существующих чанков не меняются** (passthrough остаётся byte-exact), маркер
дописывается только В КОНЦЕ, только один раз, и никогда не дублируется, если апстрим прислал свой.
## 3. Правки в коде
Файл **один**: `src/commonMain/kotlin/pw/binom/llmproxy/Main.kt`.
### 3.1. Константы и хелперы (добавить рядом с блоком think-хелперов)
```kotlin
/** Терминальный маркер SSE, который клиенты (в т.ч. Bifrost) ждут как конец потока. */
private const val SSE_DONE_MARKER = "data: [DONE]\n\n"
/** Размер скользящего окна детектора: маркер может разрезаться границей чтения. */
private const val STREAM_SCAN_WINDOW = 128
/**
* Есть ли в тексте `finish_reason` со значением, отличным от `null`. Экранированные
* вхождения (`\"finish_reason\\\":\\\"stop\\\"` внутри содержимого ответа) не считаются:
* смотрим только на неэкранированную кавычку.
*/
internal fun hasNonNullFinishReasonText(text: String): Boolean {
var from = 0
while (true) {
val i = text.indexOf("finish_reason", from)
if (i < 0) return false
val quote = i - 1
val escaped = quote >= 0 && text[quote] == '"' && quote > 0 && text[quote - 1] == '\\'
if (!escaped) {
val rest = text.substring(i + "finish_reason".length)
.dropWhile { it == ' ' || it == ':' || it == '"' }
if (!rest.startsWith("null")) return true
}
from = i + 1
}
}
/** Структурная проверка чанка: в `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()
} ?: false
/**
* Детектор контракта конца SSE для сырого passthrough-потока (`think_tags: off`).
* Копит скользящее окно последних байт, чтобы маркер `finish_reason`/`[DONE]`,
* разрезанный границей чтения, всё равно был распознан, и запоминает два факта:
* видели ли непустой `finish_reason` и видели ли `[DONE]`.
*/
internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDOW) {
var sawFinishReason = false
private set
var sawDone = false
private set
private var window: String = ""
fun feed(bytes: ByteArray, from: Int = 0, length: Int = bytes.size) {
if (length <= 0) return
val scan = window + bytes.decodeToString(from, from + length)
if (!sawDone && scan.contains("[DONE]")) sawDone = true
if (!sawFinishReason && hasNonNullFinishReasonText(scan)) sawFinishReason = true
window = if (scan.length > windowSize) scan.substring(scan.length - windowSize) else scan
}
/**
* Маркер нужен, только если апстрим закрыл поток после штатного `finish_reason`,
* но сам `[DONE]` не прислал. При обрыве БЕЗ `finish_reason` маркер НЕ дописываем.
*/
val missingDoneMarker: Boolean get() = sawFinishReason && !sawDone
}
/**
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
*/
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
val detector = StreamEndDetector()
val buf = ByteArray(8192)
while (true) {
val n = source.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
detector.feed(buf, 0, n)
}
}
if (detector.missingDoneMarker) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
```
Импорты под новые типы (`ByteReadChannel`, `ByteWriteChannel`, `readAvailable`, `writeFully`, `flush`,
`jsonArray`, `jsonObject`, `JsonPrimitive`) — проверить, что они уже есть в файле; добавлять только недостающие.
### 3.2. Правка ветки `think_tags: off` в `handleChat`
Было (внутри `call.respondBytesWriter(...)`):
```kotlin
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough как раньше.
val buf = ByteArray(8192)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
}
}
} else {
streamSseWithThinkTags(ch, thinkMode)
}
```
Стало:
```kotlin
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
ch.streamRawWithDoneContract(this)
} else {
streamSseWithThinkTags(ch, thinkMode)
}
```
(если `this` в этом контексте — не `ByteWriteChannel`, использовать тот ресивер, который уже использовался
существующими вызовами `writeFully`/`flush` в этой лямбде; компилятор — источник истины).
### 3.3. Правка `streamSseWithThinkTags`
Добавить два флага и запись фактов:
```kotlin
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break
when {
line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim()
if (payload == "[DONE]") {
sawDone = true
emitUtf8(SSE_DONE_MARKER)
} else {
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
if (obj != null && hasFinishReason(obj)) sawFinishReason = true
val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) }
if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n")
}
}
else -> emitUtf8("$line\n")
}
flush()
}
// ... существующий блок сброса удержанных хвостов сплиттеров — БЕЗ изменений ...
// НОВОЕ: после хвостов дописать маркер, если апстрим его не прислал, а finish_reason был.
if (sawFinishReason && !sawDone) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
```
Порядок важен: маркер дописывается **после** блока сброса хвостов (сначала текст, потом терминатор).
Остальную логику функции (`transformThinkChunk`, `emitUtf8("$line\n")` для служебных строк) не менять.
## 4. Тесты (новый файл)
`src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt` — по образцу существующего
`ThinkTagStreamTest` (тот же приём: `ByteChannel(autoFlush = true)`, запись входа, `close(null)`,
чтение выхода в буфер; для сырого потока вызывать `streamRawWithDoneContract`).
Обязательные случаи (имена — как ниже, чтобы можно было ссылаться в отчёте):
| Тест | Вход | Ожидание |
| --- | --- | --- |
| `textHelperSeesFinishReason` | строка `{"choices":[{"finish_reason":"stop"}]}` | `hasNonNullFinishReasonText` = true |
| `textHelperIgnoresNullFinishReason` | `{"choices":[{"finish_reason":null}]}` | false |
| `textHelperIgnoresEscapedOccurrenceInContent` | `{"delta":{"content":"\"finish_reason\":\"stop\""}}` | false |
| `detectorSeesFinishReasonSplitAcrossReads` | два `feed`: `...finish_rea` + `son":"stop"...` | `sawFinishReason` = true, `missingDoneMarker` = true |
| `detectorSeesDoneSplitAcrossReads` | два `feed`: `data: [DO` + `NE]` | `sawDone` = true, `missingDoneMarker` = false |
| `detectorSeesDoneAndFinishReason` | `finish_reason` + `[DONE]` в одном feed | `missingDoneMarker` = false |
| `detectorTruncatedStreamNeedsNoMarker` | `feed` с контентом без `finish_reason` | `missingDoneMarker` = false |
| `rawStreamIsByteExactAndAppendsDone` | SSE с `finish_reason`, без `[DONE]` | выход = вход + `"data: [DONE]\n\n"` (байт-в-байт префикс), маркер ровно один |
| `rawStreamTruncatedPassesThroughUntouched` | SSE без `finish_reason` и без `[DONE]` | выход = вход (ни одного `[DONE]`) |
| `rawStreamKeepsUpstreamDoneMarker` | SSE с `finish_reason` и своим `[DONE]` | `[DONE]` ровно один (не дублируется) |
| `thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker` | `streamSseWithThinkTags(..., "split")`, чанк с `finish_reason:"stop"`, без `[DONE]` | выход заканчивается `data: [DONE]\n\n`, ровно один маркер |
| `thinkStreamTruncatedNeedsNoMarker` | mode `split`, чанк с контентом, без `finish_reason` | в выходе нет `[DONE]` |
| `thinkStreamKeepsSingleDone` | mode `split`, `finish_reason` + `data: [DONE]` | ровно один `[DONE]` |
## 5. Мутационная приёмка (по TESTING.md)
После того как тесты зелёные, проверить их силу — по одному изменению, каждый раз с
`./gradlew cleanJvmTest jvmTest`, затем `git checkout --` на файл:
1. `missingDoneMarker` → всегда `false` (детектор сырого пути) — должны упасть тесты сырого пути
(`rawStreamIsByteExactAndAppendsDone`) и тест разрезанного маркера;
2. `if (sawFinishReason && !sawDone)` → `if (!sawDone)` в `streamSseWithThinkTags` — должен упасть
`thinkStreamTruncatedNeedsNoMarker` (это защита от «обрыв выглядит как успех»);
3. `if (sawFinishReason && !sawDone)` → `if (false)` — должен упасть
`thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker`.
Тест, который при такой поломке остаётся зелёным, считается пустым.
## 6. Команды проверки
```bash
./gradlew jvmTest # все тесты зелёные; сейчас в проекте 66 тестов, стало больше
./gradlew cleanJvmTest jvmTest # прогон с нуля (обязателен для мутационной приёмки)
./gradlew clean jvmTest fatJar # полная сборка
```
Числа тестов смотреть в `build/test-results/jvmTest/*.xml` (атрибуты `tests` / `failures` / `errors`),
а не по строкам в консоли.
## 7. Границы работ
- Трогать **только**: `src/commonMain/kotlin/pw/binom/llmproxy/Main.kt`,
`src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt`, `TESTING.md` (добавить строку про
новый тест-файл и пункты мутационной приёмки).
- **НЕ** менять `ThinkTagSplitter.kt`, non-stream путь, `transformThinkMessage`, `transformThinkChunk`,
фрейминг существующих чанков, `README.md`/`CONFIG.md`.
- **НЕ** добавлять зависимости, не менять версии в `build.gradle.kts`, не менять `config.yaml`.
- **НЕ** коммитить и не пушить.
- Внешние сети/апстримы не дёргать — только локальные тесты.
## 8. Вне области (и почему)
- **Обрыв без `finish_reason` остаётся без маркера** — так усечённый стрим не выглядит успешным
(см. п. 2). Да, в этом случае Bifrost снова будет ждать закрытия сокета; это осознанный выбор
в пользу честности контракта, а не «ускорения любой ценой».
- **Non-stream путь** не трогаем: там тула сама собирает ответ из чанков (`rebuildFromChunks`) и
маркер не нужен.
- **Bifrost и сам MiniMax** не трогаем: фикс на нашем хопе лечит всех клиентов сразу.
+1
View File
@@ -19,6 +19,7 @@ kotlinx-serialization-json = { module = "org.jetbrains.kotlinx:kotlinx-serializa
yamlkt = { module = "net.mamoe.yamlkt:yamlkt", version.ref = "yamlkt" }
kotlinx-coroutines-core = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-core", version.ref = "coroutines" }
kotlinx-datetime = { module = "org.jetbrains.kotlinx:kotlinx-datetime", version.ref = "datetime" }
kotlinx-coroutines-test = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-test", version.ref = "coroutines" }
kotlinx-io-core = { module = "org.jetbrains.kotlinx:kotlinx-io-core", version.ref = "kotlinxIo" }
logback-classic = { module = "ch.qos.logback:logback-classic", version.ref = "logback" }
kotlin-test = { module = "org.jetbrains.kotlin:kotlin-test" }
@@ -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()
}
}
+303 -46
View File
@@ -46,6 +46,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 +90,19 @@ 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] server: bind host=${config.server.host} port=${config.server.port}"
}
@@ -101,7 +113,7 @@ fun main() {
val http = createHttpClient()
val sessions = SessionRegistry()
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>,
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 +154,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 +188,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 +199,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,13 +252,29 @@ 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)
@@ -242,20 +282,13 @@ private suspend fun handleChat(
call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>()
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough как раньше.
val buf = ByteArray(8192)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
}
}
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
streamRawWithDoneContract(ch)
} else {
streamSseWithThinkTags(ch, thinkMode)
}
}
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"
@@ -279,6 +312,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"
@@ -289,10 +323,16 @@ private suspend fun handleChat(
if (failover) continue
if (responded) return
} 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 {
@@ -302,10 +342,20 @@ private suspend fun handleChat(
if (anyClaimed) {
call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
} else {
val retryIn = pool.map { backoff.remainingSeconds(it) }.filter { it > 0 }.minOrNull()
if (retryIn != null) {
call.response.headers.append("Retry-After", retryIn.toString())
call.respondText(
errorJson("all upstreams cooling, retry in ${retryIn}s"),
ContentType.Application.Json,
HttpStatusCode.ServiceUnavailable,
)
} else {
call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
}
}
}
private fun errorJson(msg: String): String =
"""{"error":{"message":"$msg"}}"""
@@ -487,15 +537,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)}" }
@@ -546,27 +609,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 }
}
}
@@ -618,12 +681,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)
@@ -644,6 +707,144 @@ 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`
* добавляем это поле, ЕСЛИ его там ещё нет. Текст берём из `reasoning`
* (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`).
* Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста
* пишем пустую строку, только если `emptyOk`.
* Тело возвращается без изменений (тот же объект), если менять нечего.
*/
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 toolCalls = msg["tool_calls"] as? JsonArray
if (toolCalls == null || toolCalls.isEmpty()) return@map el
val existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (existing != null && existing.isNotEmpty()) return@map el
val text = reasoningTextOf(msg)
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"
/** Размер скользящего окна детектора: маркер может разрезаться границей чтения. */
private const val STREAM_SCAN_WINDOW = 128
/**
* Есть ли в тексте `finish_reason` со значением, отличным от `null`. Экранированные
* вхождения (`\"finish_reason\\\":\\\"stop\\\"` внутри содержимого ответа) не считаются:
* смотрим только на неэкранированную кавычку.
*/
internal fun hasNonNullFinishReasonText(text: String): Boolean {
var from = 0
while (true) {
val i = text.indexOf("finish_reason", from)
if (i < 0) return false
val quote = i - 1
val escaped = quote >= 0 && text[quote] == '"' && quote > 0 && text[quote - 1] == '\\'
if (!escaped) {
val rest = text.substring(i + "finish_reason".length)
.dropWhile { it == ' ' || it == ':' || it == '"' }
if (!rest.startsWith("null")) return true
}
from = i + 1
}
}
/** Структурная проверка чанка: в `choices[*].finish_reason` есть непустая строка. */
internal fun hasFinishReason(obj: JsonObject): Boolean =
(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
/**
* Детектор контракта конца SSE для сырого passthrough-потока (`think_tags: off`).
* Копит скользящее окно последних байт, чтобы маркер `finish_reason`/`[DONE]`,
* разрезанный границей чтения, всё равно был распознан, и запоминает два факта:
* видели ли непустой `finish_reason` и видели ли `[DONE]`.
*/
internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDOW) {
var sawFinishReason = false
private set
var sawDone = false
private set
private var window: String = ""
fun feed(bytes: ByteArray, from: Int = 0, length: Int = bytes.size) {
if (length <= 0) return
val scan = window + bytes.decodeToString(from, from + length)
if (!sawDone && scan.contains("[DONE]")) sawDone = true
if (!sawFinishReason && hasNonNullFinishReasonText(scan)) sawFinishReason = true
window = if (scan.length > windowSize) scan.substring(scan.length - windowSize) else scan
}
/**
* Маркер нужен, только если апстрим закрыл поток после штатного `finish_reason`,
* но сам `[DONE]` не прислал. При обрыве БЕЗ `finish_reason` маркер НЕ дописываем.
*/
val missingDoneMarker: Boolean get() = sawFinishReason && !sawDone
}
/**
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
*/
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
val detector = StreamEndDetector()
val buf = ByteArray(8192)
while (true) {
val n = source.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
detector.feed(buf, 0, n)
}
}
if (detector.missingDoneMarker) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
/**
* Построчный разбор SSE-стрима с рассечением think-тегов. Строки, не начинающиеся
* с `data:`, и `data: [DONE]` уходят клиенту без изменений (с `\n`). Прочие
@@ -652,18 +853,22 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
*/
private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break
when {
line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim()
if (payload == "[DONE]") {
emitUtf8("data: [DONE]\n\n")
sawDone = true
emitUtf8(SSE_DONE_MARKER)
} else {
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
if (obj != null && hasFinishReason(obj)) sawFinishReason = true
val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) }
if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n")
}
@@ -699,6 +904,12 @@ private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChan
flush()
}
}
// После сброса хвостов дописываем маркер, только если апстрим его не прислал,
// а `finish_reason` в потоке был (при обрыве без него маркер НЕ дописываем).
if (sawFinishReason && !sawDone) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
/**
@@ -710,20 +921,20 @@ private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChan
* Чанк без `choices` или без строкового `delta.content` не меняется —
* возвращается null (отдать исходную строку как есть).
*/
private fun transformThinkChunk(
internal fun transformThinkChunk(
obj: JsonObject,
splitters: MutableMap<Int, ThinkTagSplitter>,
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
@@ -751,6 +962,9 @@ 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,
)
data class UpstreamConf(
@@ -760,6 +974,7 @@ data class UpstreamConf(
val max_concurrency: Int? = null,
val patch: JsonObject? = null,
val think_tags: String? = null,
val backoff: Duration? = null,
)
data class ModelConf(
@@ -815,6 +1030,9 @@ 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"),
)
}
@@ -827,6 +1045,7 @@ 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"),
)
}
@@ -852,5 +1071,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,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.assertTrue
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 +15,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 +225,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
@@ -533,7 +536,7 @@ class ConfigLogicTest {
// значение у апстрима — берётся оно, провайдер игнорируется
assertEquals("split", effectiveThinkTags(upSplit, null))
assertEquals("split", effectiveThinkTags(upSplit, prov("true")))
assertEquals("split", effectiveThinkTags(upSplit, prov("strip")))
// апстрим null — берётся провайдерский
assertEquals("split", effectiveThinkTags(upNull, prov("split")))
@@ -546,4 +549,41 @@ 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)
}
}
@@ -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,165 @@
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 assistantWithoutToolCallsIsUntouched() {
// Без tool_calls assistant-сообщение не трогаем (тело идентично).
val m = """{"role":"assistant","reasoning":"R","content":"hi"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
}
@Test
fun assistantWithEmptyToolCallsIsUntouched() {
// Пустой массив tool_calls = не трогаем (тело идентично).
val m = """{"role":"assistant","tool_calls":[],"reasoning":"R"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
}
@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"))
}
}
@@ -0,0 +1,176 @@
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.assertFalse
import kotlin.test.assertTrue
import kotlinx.coroutines.test.runTest
class StreamDoneContractTest {
/** Терминальный маркер, который должен оказаться (или не оказаться) в конце вывода. */
private val done = "data: [DONE]\n\n"
/** Прогнать сырой passthrough-поток через боевую обвязку и вернуть записанные байты как строку. */
private suspend fun runRaw(input: String): String {
val src = ByteChannel(autoFlush = true)
src.writeFully(input.encodeToByteArray())
src.close(null)
val out = ByteChannel(autoFlush = true)
out.streamRawWithDoneContract(src)
out.close(null)
return readAll(out)
}
/** Прогнать SSE через think-обвязку и вернуть записанные байты как строку. */
private suspend fun runThink(input: String, mode: String): String {
val src = ByteChannel(autoFlush = true)
src.writeFully(input.encodeToByteArray())
src.close(null)
val out = ByteChannel(autoFlush = true)
out.streamSseWithThinkTags(src, mode)
out.close(null)
return readAll(out)
}
private suspend fun readAll(ch: ByteChannel): String {
val sb = StringBuilder()
val buf = ByteArray(512)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) sb.append(buf.decodeToString(0, n))
}
return sb.toString()
}
private fun countOccurrences(text: String, needle: String): Int {
var count = 0
var from = 0
while (true) {
val i = text.indexOf(needle, from)
if (i < 0) return count
count++
from = i + needle.length
}
}
@Test
fun textHelperSeesFinishReason() {
assertTrue(hasNonNullFinishReasonText("""{"choices":[{"finish_reason":"stop"}]}"""))
}
@Test
fun textHelperIgnoresNullFinishReason() {
assertFalse(hasNonNullFinishReasonText("""{"choices":[{"finish_reason":null}]}"""))
}
@Test
fun textHelperIgnoresEscapedOccurrenceInContent() {
// Экранированное вхождение внутри content — это не поле finish_reason.
assertFalse(hasNonNullFinishReasonText("""{"delta":{"content":"\"finish_reason\":\"stop\""}}"""))
}
@Test
fun detectorSeesFinishReasonSplitAcrossReads() {
// Маркер finish_reason разрезан границей чтения: детектор обязан его собрать.
val detector = StreamEndDetector()
detector.feed("""{"choices":[{"finish_rea""".encodeToByteArray())
detector.feed("""son":"stop"}]}""".encodeToByteArray())
assertTrue(detector.sawFinishReason)
assertTrue(detector.missingDoneMarker)
}
@Test
fun detectorSeesDoneSplitAcrossReads() {
// [DONE] разрезан границей чтения: детектор видит его и маркер не нужен.
val detector = StreamEndDetector()
detector.feed("data: [DO".encodeToByteArray())
detector.feed("NE]".encodeToByteArray())
assertTrue(detector.sawDone)
assertFalse(detector.missingDoneMarker)
}
@Test
fun detectorSeesDoneAndFinishReason() {
// В одном feed и finish_reason, и [DONE] — дописывать ничего не нужно.
val detector = StreamEndDetector()
detector.feed("""{"choices":[{"finish_reason":"stop"}]} data: [DONE]""".encodeToByteArray())
assertFalse(detector.missingDoneMarker)
}
@Test
fun detectorTruncatedStreamNeedsNoMarker() {
// Обрыв без finish_reason: маркер не дописываем (усечённый стрим честно остаётся усечённым).
val detector = StreamEndDetector()
detector.feed("""{"choices":[{"delta":{"content":"hi"}}]}""".encodeToByteArray())
assertFalse(detector.missingDoneMarker)
}
@Test
fun rawStreamIsByteExactAndAppendsDone() = runTest {
// Сырой поток с finish_reason и без [DONE]: байты входа не меняются, маркер дописывается ровно один.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"},\"finish_reason\":\"stop\"}]}\n\n"
val out = runRaw(input)
assertEquals(input + done, out, "выход должен быть вход + один маркер")
assertEquals(1, countOccurrences(out, done), "маркер должен быть ровно один")
}
@Test
fun rawStreamTruncatedPassesThroughUntouched() = runTest {
// Обрыв без finish_reason: ни одного [DONE], вход не меняется.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
val out = runRaw(input)
assertEquals(input, out, "усечённый поток должен уйти как есть")
assertEquals(0, countOccurrences(out, "[DONE]"))
}
@Test
fun rawStreamKeepsUpstreamDoneMarker() = runTest {
// Апстрим прислал свой [DONE]: второй маркер не дописывается.
val input = "data: {\"choices\":[{\"finish_reason\":\"stop\"}]}\n\n" + done
val out = runRaw(input)
assertEquals(1, countOccurrences(out, done), "маркер не должен дублироваться")
}
@Test
fun thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker() = runTest {
// think-обвязка: finish_reason есть, [DONE] нет — маркер дописывается ровно один.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"},\"finish_reason\":\"stop\"}]}\n\n"
val out = runThink(input, "split")
assertTrue(out.endsWith(done), "выход должен заканчиваться маркером: $out")
assertEquals(1, countOccurrences(out, done), "маркер должен быть ровно один")
}
@Test
fun thinkStreamTruncatedNeedsNoMarker() = runTest {
// think-обвязка: обрыв без finish_reason — маркер не дописываем.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
val out = runThink(input, "split")
assertEquals(0, countOccurrences(out, "[DONE]"), "при обрыве маркера быть не должно: $out")
}
@Test
fun thinkStreamKeepsSingleDone() = runTest {
// think-обвязка: finish_reason и собственный [DONE] — маркер остаётся ровно один.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"},\"finish_reason\":\"stop\"}]}\n\n" + done
val out = runThink(input, "split")
assertEquals(1, countOccurrences(out, done), "маркер не должен дублироваться")
}
}
@@ -0,0 +1,179 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class ThinkTagChunkTest {
// Литералы тегов OPEN/CLOSE из ThinkTagSplitter собираем из символьных
// кусочков, чтобы не записывать тег единой строкой в исходнике.
private val openTag = '<' + "think" + '>'
private val closeTag = '<' + "/think" + '>'
// Маркер рассуждения (5 букв) собираем из Unicode-кодов, без единой строки.
private val reasoning = listOf(0x0420, 0x0410, 0x0417, 0x0423, 0x041C)
.map { Char(it) }
.joinToString("")
private fun deltaOf(result: String, index: Int): JsonObject =
Json.parseToJsonElement(result).jsonObject
.get("choices")!!.jsonArray[index].jsonObject
.get("delta")!!.jsonObject
@Test
fun chunkSplitHoldsTailAcrossChunks() {
// Хвост открывающего тега, разрезанный на границе чанков, переживает
// два отдельных вызова через общий splitters (getOrPut по index).
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
// Чанк 1: текст + обрезанный префикс открывающего тега (не завершает тег).
val prefix = openTag.dropLast(1)
val obj1 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"текст$prefix"}}]}""",
).jsonObject
val r1 = transformThinkChunk(obj1, splitters, "split", true)
assertNotNull(r1)
val d1 = deltaOf(r1, 0)
assertEquals("текст", d1["content"]!!.jsonPrimitive.content)
assertFalse(d1.containsKey("reasoning_content"))
// Чанк 2 (тот же splitters): остаток открывающего тега + маркер.
val rest = openTag.last()
val obj2 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"$rest$reasoning"}}]}""",
).jsonObject
val r2 = transformThinkChunk(obj2, splitters, "split", true)
assertNotNull(r2)
val d2 = deltaOf(r2, 0)
assertFalse(d2.containsKey("content"))
assertEquals(reasoning, d2["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun chunkStripRemovesReasoning() {
// Режим strip: блок целиком (теги + рассуждения) вырезается,
// content склеивается в «AB», ключа reasoning_content нет.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val c = "A" + openTag + reasoning + closeTag + "B"
val obj = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"$c"}}]}""",
).jsonObject
val r = transformThinkChunk(obj, splitters, "strip", false)
assertNotNull(r)
val d = deltaOf(r, 0)
assertEquals("AB", d["content"]!!.jsonPrimitive.content)
assertFalse(d.containsKey("reasoning_content"))
}
@Test
fun chunkDropsEmptyContentFieldWhenAllGoesToReasoning() {
// Режим split: весь контент чанка — think-блок, поэтому content пуст
// и ключ убирается из delta, а рассуждения уходят в reasoning_content.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val c = openTag + reasoning + closeTag
val obj = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"$c"}}]}""",
).jsonObject
val r = transformThinkChunk(obj, splitters, "split", true)
assertNotNull(r)
val d = deltaOf(r, 0)
assertFalse(d.containsKey("content"))
assertEquals(reasoning, d["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun chunkReturnsNullWhenNoChoices() {
// Нет ключа choices — чанк не трогаем, отдаём исходную строку (null).
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val obj = Json.parseToJsonElement("""{"model":"m"}""").jsonObject
assertNull(transformThinkChunk(obj, splitters, "split", true))
}
@Test
fun chunkReturnsNullWhenDeltaContentNotAString() {
// content = null (JsonNull) и content = массив частей — не строка,
// значит choice не трогаем, весь чанк не меняется -> null.
val s1 = mutableMapOf<Int, ThinkTagSplitter>()
val o1 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":null}}]}""",
).jsonObject
assertNull(transformThinkChunk(o1, s1, "split", true))
val s2 = mutableMapOf<Int, ThinkTagSplitter>()
val o2 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":[{"type":"text","text":"hi"}]}}]}""",
).jsonObject
assertNull(transformThinkChunk(o2, s2, "split", true))
}
@Test
fun chunkSeparateSplittersPerIndex() {
// Один объект splitters на все три вызова: по каждому index держим
// своего сплиттера, поэтому обрезанные хвосты по index не путаются.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val prefix = openTag.dropLast(1)
val rest = openTag.last()
// Вызов 1: два choice — index 0 «A»+префикс тега, index 1 «B»+тот же префикс.
val c1 = """{"choices":[{"index":0,"delta":{"content":"A$prefix"}},{"index":1,"delta":{"content":"B$prefix"}}]}"""
val r1 = transformThinkChunk(Json.parseToJsonElement(c1).jsonObject, splitters, "split", true)
assertNotNull(r1)
// Вызов 2: index 1 получает остаток тега + маркер → рассуждения у index 1.
val c2 = """{"choices":[{"index":1,"delta":{"content":"$rest${reasoning}1"}}]}"""
val r2 = transformThinkChunk(Json.parseToJsonElement(c2).jsonObject, splitters, "split", true)
assertNotNull(r2)
val d2 = deltaOf(r2, 0)
assertEquals(reasoning + "1", d2["reasoning_content"]!!.jsonPrimitive.content)
// Вызов 3: index 0 получает остаток тега + маркер → рассуждения у index 0.
val c3 = """{"choices":[{"index":0,"delta":{"content":"$rest${reasoning}0"}}]}"""
val r3 = transformThinkChunk(Json.parseToJsonElement(c3).jsonObject, splitters, "split", true)
assertNotNull(r3)
val d3 = deltaOf(r3, 0)
assertEquals(reasoning + "0", d3["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun chunkKeepsUntouchedChoiceAndFinishReason() {
// Первый choice целиком идёт через split; второй без ключа content
// (в нём только finish_reason) — остаётся нетронутым, как и был.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val c = "A" + openTag + reasoning + closeTag + "B"
val c1 = """{"choices":[{"index":0,"delta":{"content":"$c"}},{"index":1,"delta":{"finish_reason":"stop"}}]}"""
val r = transformThinkChunk(Json.parseToJsonElement(c1).jsonObject, splitters, "split", true)
assertNotNull(r)
val d0 = deltaOf(r, 0)
assertEquals("AB", d0["content"]!!.jsonPrimitive.content)
assertEquals(reasoning, d0["reasoning_content"]!!.jsonPrimitive.content)
val d1 = deltaOf(r, 1)
assertEquals("stop", d1["finish_reason"]!!.jsonPrimitive.content)
assertFalse(d1.containsKey("content"))
}
@Test
fun chunkIndexDefaultsToZeroWhenAbsent() {
// Ключ index отсутствует — оба раза используем сплиттер под индекс 0,
// поэтому хвост первого чанка доживал до рассуждений второго.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val prefix = openTag.dropLast(1)
val rest = openTag.last()
val c1 = """{"choices":[{"delta":{"content":"A$prefix"}}]}"""
val r1 = transformThinkChunk(Json.parseToJsonElement(c1).jsonObject, splitters, "split", true)
assertNotNull(r1)
val c2 = """{"choices":[{"delta":{"content":"$rest$reasoning"}}]}"""
val r2 = transformThinkChunk(Json.parseToJsonElement(c2).jsonObject, splitters, "split", true)
assertNotNull(r2)
val d2 = deltaOf(r2, 0)
assertEquals(reasoning, d2["reasoning_content"]!!.jsonPrimitive.content)
}
}
@@ -0,0 +1,179 @@
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.assertTrue
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class ThinkTagStreamTest {
/** Прогнать SSE-текст через боевую обвязку и вернуть то, что она записала. */
private suspend fun runStream(input: String, mode: String): String {
val src = ByteChannel(autoFlush = true)
src.writeFully(input.encodeToByteArray())
src.close(null)
val out = ByteChannel(autoFlush = true)
out.streamSseWithThinkTags(src, mode)
out.close(null)
val sb = StringBuilder()
val buf = ByteArray(512)
while (true) {
val n = out.readAvailable(buf)
if (n == -1) break
if (n > 0) sb.append(buf.decodeToString(0, n))
}
return sb.toString()
}
// Литералы тегов OPEN/CLOSE из ThinkTagSplitter собираем из символьных
// кусочков, чтобы не записывать тег единой строкой в исходнике.
private val openTag = '<' + "think" + '>'
private val closeTag = '<' + "/think" + '>'
// Маркер рассуждения (5 букв) собираем из Unicode-кодов, без единой строки.
private val reasoning = listOf(0x0420, 0x0410, 0x0417, 0x0423, 0x041C)
.map { Char(it) }
.joinToString("")
/** JSON-нагрузки из data-строк вывода (строки, начинающиеся с «data:»), без служебного [DONE]. */
private fun dataPayloads(sse: String): List<String> =
sse.lines()
.filter { it.startsWith("data:") && it.removePrefix("data:").trim() != "[DONE]" }
.map { it.removePrefix("data:").trim() }
private fun deltaOf(result: String, index: Int): JsonObject =
Json.parseToJsonElement(result).jsonObject
.get("choices")!!.jsonArray[index].jsonObject
.get("delta")!!.jsonObject
@Test
fun streamDoneAndServiceLinesPassThrough() = runTest {
// Служебные строки SSE и маркер конца потока должны дойти до клиента без изменений.
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n")
append("\n")
append("data: [DONE]\n")
append("\n")
append("event: ping\n")
append("\n")
append(": keep-alive\n")
append("\n")
}
val out = runStream(input, "split")
assertTrue(out.contains("data: [DONE]"), "маркер конца потока потерян: $out")
assertTrue(out.contains("event: ping"), "служебная строка event потеряна: $out")
assertTrue(out.contains(": keep-alive"), "строка-комментарий потеряна: $out")
assertTrue(out.contains("hi"), "текстовый чанк потерян: $out")
}
@Test
fun streamSplitsReasoningAcrossDataChunks() = runTest {
// Хвост открывающего тега разрезан на границе двух data-событий: первое
// отдаёт content «A» и хвост держит, второе завершает тег и блок.
val prefix = openTag.take(3)
val rest = openTag.drop(3)
val content2 = rest + reasoning + closeTag + "B"
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"A$prefix\"}}]}\n")
append("\n")
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"$content2\"}}]}\n")
append("\n")
}
val out = runStream(input, "split")
val payloads = dataPayloads(out)
val last = deltaOf(payloads.last(), 0)
assertEquals("B", last["content"]!!.jsonPrimitive.content)
assertEquals(reasoning, last["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun streamUnclosedBlockCarriesReasoningInSameChunk() = runTest {
// Незакрытый think-блок в конце потока: reasoning отдаётся сразу в том же
// чанке, где пришёл, финиш-чанка не появляется — удержанного хвоста нет.
val content = "A" + openTag + "МЫСЛИ"
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"$content\"}}]}\n")
append("\n")
}
val payloads = dataPayloads(runStream(input, "split"))
assertEquals(1, payloads.size, "финиш-чанка быть не должно: $payloads")
val delta = deltaOf(payloads.single(), 0)
assertEquals("A", delta["content"]!!.jsonPrimitive.content)
assertEquals("МЫСЛИ", delta["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun streamFlushesHeldTailOnFinish() = runTest {
// Поток обрывается на неполном префиксе открывающего тега: удержанный
// хвост не теряется и сбрасывается финиш-чанком в конце.
val tail = openTag.take(3)
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"A$tail\"}}]}\n")
append("\n")
}
val payloads = dataPayloads(runStream(input, "split"))
assertEquals(2, payloads.size, "финиш-чанк с удержанным хвостом потерян: $payloads")
val flushed = deltaOf(payloads.last(), 0)
assertEquals(tail, flushed["content"]!!.jsonPrimitive.content)
}
@Test
fun streamInvalidJsonPassesThroughVerbatim() = runTest {
// data-строка, которая не является JSON, доходит до клиента без изменений.
val input = buildString {
append("data: {это не json\n")
append("\n")
}
val out = runStream(input, "split")
assertTrue(out.contains("data: {это не json"), "битая строка не дошла как есть: $out")
}
@Test
fun streamChunkWithoutChoicesPassesThrough() = runTest {
// Чанк с пустыми choices не меняется и уходит клиенту исходной строкой.
val input = buildString {
append("data: {\"usage\":{\"total_tokens\":5},\"choices\":[]}\n")
append("\n")
}
val out = runStream(input, "split")
assertTrue(out.contains("\"total_tokens\":5"), "чанк с usage изменился: $out")
}
@Test
fun streamModeStripEmitsNoReasoning() = runTest {
// Режим strip: текст внутри think-тегов выбрасывается, поля
// reasoning_content нет, обычный текст остаётся.
val content = "A" + openTag + reasoning + closeTag + "B"
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"$content\"}}]}\n")
append("\n")
}
val out = runStream(input, "strip")
val payloads = dataPayloads(out)
assertTrue(!out.contains("reasoning_content"), "в режиме strip не должно быть reasoning_content: $out")
assertEquals(1, payloads.size, "ожидался ровно один чанк: $payloads")
assertEquals("AB", deltaOf(payloads.single(), 0)["content"]!!.jsonPrimitive.content, "теги и рассуждение должны быть вырезаны: $payloads")
}
}
@@ -0,0 +1,174 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class ThinkTagTransformTest {
// Точные литералы тегов из ThinkTagSplitter (OPEN/CLOSE): собираем из
// отдельных символов, чтобы не писать тег единой строкой в исходнике.
private val openTag = '<' + "think" + '>'
private val closeTag = '<' + "/think" + '>'
private fun messageOf(result: String, index: Int): JsonObject =
Json.parseToJsonElement(result).jsonObject
.get("choices")!!.jsonArray[index].jsonObject
.get("message")!!.jsonObject
@Test
fun splitSingleBlockMovesReasoningAndKeepsOtherFields() {
// Проверяем: в режиме split содержимое блока уходит в reasoning_content,
// в content остаётся «AB», а служебные поля (model, usage, …) не тронуты.
val c = "A" + openTag + "РАЗУМ" + closeTag + "B"
val obj = Json.parseToJsonElement(
"""
{"id":"c1","object":"chat.completion","created":123,"model":"m1",
"usage":{"prompt_tokens":1,"completion_tokens":2,"total_tokens":3},
"system_fingerprint":"sf1",
"choices":[{"index":0,"message":{"role":"assistant","content":"$c"},"finish_reason":"stop"}]}
""".trimIndent(),
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val out = Json.parseToJsonElement(result).jsonObject
val msg = out.get("choices")!!.jsonArray[0].jsonObject.get("message")!!.jsonObject
assertEquals("AB", msg["content"]!!.jsonPrimitive.content)
assertEquals("РАЗУМ", msg["reasoning_content"]!!.jsonPrimitive.content)
assertEquals("m1", out["model"]!!.jsonPrimitive.content)
assertEquals(3, out["usage"]!!.jsonObject["total_tokens"]!!.jsonPrimitive.content.toInt())
}
@Test
fun splitTwoBlocksConcatenateInOrder() {
// Проверяем порядок конкатенации reasoning, когда два think-блока идут подряд.
val c = openTag + "ПЕРВЫЙ" + closeTag + "X" + openTag + "ВТОРОЙ" + closeTag
val obj = Json.parseToJsonElement(
"""{"choices":[{"index":0,"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("X", msg["content"]!!.jsonPrimitive.content)
assertEquals("ПЕРВЫЙВТОРОЙ", msg["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun splitAppendsToExistingReasoningContent() {
// Проверяем, что прежнее reasoning_content не теряется, а новое дописывается.
val c = openTag + "NEW" + closeTag
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c","reasoning_content":"OLD"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("OLDNEW", msg["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun splitUnclosedTagGoesToReasoning() {
// Проверяем, что незакрытый открывающий тег уводит хвост целиком в reasoning.
val c = "текст" + openTag + "мысли"
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("текст", msg["content"]!!.jsonPrimitive.content)
assertEquals("мысли", msg["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun stripRemovesTagsAndLeavesNoReasoningKey() {
// Проверяем: в режиме strip теги и рассуждения вырезаются,
// а ключ reasoning_content в message отсутствует вовсе.
val c = "A" + openTag + "РАЗУМ" + closeTag + "B"
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "strip")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("AB", msg["content"]!!.jsonPrimitive.content)
assertFalse(msg.containsKey("reasoning_content"))
}
@Test
fun missingChoicesReturnsNull() {
// Проверяем: без ключа choices функция не вносит изменений (null).
val obj = Json.parseToJsonElement(
"""{"id":"c1","model":"m1"}""",
).jsonObject
assertNull(transformThinkMessage(obj, "split"))
}
@Test
fun nonStringContentReturnsNull() {
// content = null (JsonNull) — менять нечего, функция возвращает null.
val nullContent = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":null}}]}""",
).jsonObject
assertNull(transformThinkMessage(nullContent, "split"))
// content = массив частей (JsonArray) — тоже не трогаем, null.
val arrayContent = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":[{"type":"text","text":"hi"}]}}]}""",
).jsonObject
assertNull(transformThinkMessage(arrayContent, "split"))
}
@Test
fun twoChoicesEachKeepOwnThinkBlock() {
// Проверяем, что каждый choice обрабатывается независимо: у каждого свой
// вырезанный content и свой reasoning_content.
val c0 = "А" + openTag + "А-раз" + closeTag + "Б"
val c1 = "В" + openTag + "В-раз" + closeTag + "Г"
val obj = Json.parseToJsonElement(
"""
{"choices":[
{"index":0,"message":{"content":"$c0"},"finish_reason":"stop"},
{"index":1,"message":{"content":"$c1"},"finish_reason":"stop"}]}
""".trimIndent(),
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val out = Json.parseToJsonElement(result).jsonObject.get("choices")!!.jsonArray
val m0 = out[0].jsonObject.get("message")!!.jsonObject
val m1 = out[1].jsonObject.get("message")!!.jsonObject
assertEquals("АБ", m0["content"]!!.jsonPrimitive.content)
assertEquals("А-раз", m0["reasoning_content"]!!.jsonPrimitive.content)
assertEquals("ВГ", m1["content"]!!.jsonPrimitive.content)
assertEquals("В-раз", m1["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun offModePassthroughKeepsTags() {
// Проверяем passthrough: в режиме off теги остаются в content,
// reasoning_content не добавляется, а результат — не null.
val c = "A" + openTag + "РАЗУМ" + closeTag + "B"
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "off")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("A" + openTag + "РАЗУМ" + closeTag + "B", msg["content"]!!.jsonPrimitive.content)
assertFalse(msg.containsKey("reasoning_content"))
}
}