Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c6a24a94b0 | |||
| 249b83eb51 | |||
| 293964d8f1 |
@@ -85,6 +85,8 @@ providers:
|
||||
allow_fallbacks: false
|
||||
backoff: P1M # опционально; потолок экспоненциального backoff
|
||||
# на весь провайдер (ISO-8601: P1M, PT15M, P1D…)
|
||||
probe_interval: PT1S # опционально; фоновый пинг апстримов в откате
|
||||
# (как часто проверять, жив ли провайдер)
|
||||
|
||||
- id: local-llama
|
||||
url: "http://10.0.0.5:8080/v1"
|
||||
@@ -103,6 +105,8 @@ upstreams:
|
||||
ignore: [deepseek]
|
||||
backoff: PT15M # опционально; собственный backoff этой модели
|
||||
# (приоритет над backoff провайдера)
|
||||
probe_interval: PT0S # опционально; пинг отключён у этой модели
|
||||
# (приоритет над probe_interval провайдера)
|
||||
|
||||
- id: local-qwen
|
||||
provider: local-llama
|
||||
@@ -231,12 +235,15 @@ upstreams:
|
||||
### Нативное поле рассуждений (`reasoning_field`)
|
||||
|
||||
Некоторые шлюзы-апстримы (например, **Console Go / deepseek в thinking-режиме**)
|
||||
в thinking-режиме **требуют** вернуть им нативное поле рассуждений в каждом
|
||||
assistant-сообщении с непустым `tool_calls`. Если его нет — апстрим отвечает
|
||||
`400` (`The `reasoning_content` in the thinking mode must be passed back to the
|
||||
в thinking-режиме **требуют** вернуть им нативное поле рассуждений в **каждом**
|
||||
assistant-сообщении истории (не только в тех, что с `tool_calls`). Если его нет —
|
||||
апстрим отвечает `400`
|
||||
(`The `reasoning_content` in the thinking mode must be passed back to the
|
||||
API.`). Клиенты при этом рассуждения держат в своих форматах: `reasoning`
|
||||
(строка) и/или `reasoning_details` (массив `{type:"reasoning.text", text, ...}`)
|
||||
— а нативное поле `reasoning_content` в истории могут и не передавать.
|
||||
(Точно так же делает и сам opencode: для deepseek он добавляет reasoning-часть
|
||||
на каждом assistant-сообщении, даже пустую.)
|
||||
|
||||
Поля (только на уровне **провайдера**, это свойство шлюза, а не модели):
|
||||
|
||||
@@ -245,15 +252,15 @@ API.`). Клиенты при этом рассуждения держат в с
|
||||
| `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) |
|
||||
| `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе |
|
||||
|
||||
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в каждом
|
||||
assistant-сообщении с непустым `tool_calls` — берёт текст из `reasoning`
|
||||
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в
|
||||
**каждом** assistant-сообщении — берёт текст из `reasoning`
|
||||
(если это непустая строка), иначе склеивает `reasoning_details[*].text`
|
||||
(только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и
|
||||
записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее
|
||||
непустое поле не перезаписывается. Ничего при этом не убирается и не
|
||||
переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте,
|
||||
поле просто дополняется. Тронуто только assistant-сообщение с непустым
|
||||
`tool_calls`; сообщения без `tool_calls` и `user`/`tool`/`system` не меняются.
|
||||
поле просто дополняется. Тронуты только assistant-сообщения; `user`/`tool`/
|
||||
`system` не меняются.
|
||||
Если текста рассуждений нет вовсе — поле не добавляется, кроме случая
|
||||
`reasoning_empty_ok: true` (тогда пишется пустая строка `""`). Если
|
||||
`reasoning_field` не задан — тело не меняется вовсе.
|
||||
@@ -319,6 +326,131 @@ upstreams:
|
||||
При старте выводится, что настроено:
|
||||
`[llm-proxy] backoff: upstreams=routerai-gpt4o=PT15M providers=routerai=P1M`.
|
||||
|
||||
### Таймаут первого байта стрима (`first_byte_timeout`)
|
||||
|
||||
Для `stream=true` запросов прокси следит, чтобы апстрим начал отдавать
|
||||
**что-нибудь** (первый байт тела) в течение заданного окна. Если апстрим
|
||||
молчит дольше — это трактуется как его сбой: слот конкурентности
|
||||
освобождается, ошибка идёт в backoff-откат, и роутер **фейловерит** на
|
||||
следующего свободного апстрима модели (при его наличии). Это закрывает
|
||||
кейс, когда апстрим принял запрос, но «завис» и никогда не начинает
|
||||
стримить.
|
||||
|
||||
| Поле | Тип / дефолт | Значение |
|
||||
|---|---|---|
|
||||
| `providers[].first_byte_timeout` | ISO-8601-длительность / отсутствует | Окно ожидания первого байта для всех моделей провайдера |
|
||||
| `upstreams[].first_byte_timeout` | ISO-8601-длительность / отсутствует | Окно для конкретной модели (приоритет над провайдерским) |
|
||||
|
||||
**Дефолт** (если не задан нигде): `PT30S`.
|
||||
|
||||
**Приоритет — модели.** Если `upstreams[].first_byte_timeout` задан, он
|
||||
перекрывает провайдерский; если нет — берётся провайдерский; если нет и там —
|
||||
дефолт `PT30S`.
|
||||
|
||||
**Отключение:** `PT0S` на провайдере или апстриме отключает watchdog
|
||||
для этого апстрима — тогда первый байт ждём без лимита (в рамках общего
|
||||
`upstream: request_timeout`, см. ниже).
|
||||
|
||||
**Применяется только к `stream=true`.** Для `stream=false` (не-стриминговый
|
||||
ответ) действует только общий лимит `UPSTREAM_REQUEST_TIMEOUT_MS` (`5m`) —
|
||||
строки SSE там склеиваются в один ответ, и таймаут первого байта к нему
|
||||
неприменим.
|
||||
|
||||
Поведение:
|
||||
|
||||
- watchdog висит **на первом чтении** из тела стрим-ответа; как только
|
||||
пришёл хоть какой-то байт — лимит отключается, дальше стрим читается
|
||||
свободно (вплоть до общего request_timeout);
|
||||
- `data: [DONE]` приходит в теле — это уже «первый байт», таймаут не
|
||||
сработает;
|
||||
- при срабатывании — лог:
|
||||
`[llm-proxy] chat model=<m> upstream=<id> TIMEOUT: первый байт не пришёл за <N>ms → фейловер`
|
||||
(+ `(backoff: отдых …)` если задан backoff), метрика
|
||||
`llm_proxy_requests_total{result="timeout"}`;
|
||||
- закрытие апстримом пустого стрима (EOF без единого байта) — **не** таймаут:
|
||||
это валидный короткий ответ, он проксируется как обычно.
|
||||
|
||||
Значение — ISO-8601-длительность, формат как у `backoff` (`PT30S`,
|
||||
`PT1M30S`, `P1D`, …).
|
||||
|
||||
```yaml
|
||||
providers:
|
||||
- id: routerai
|
||||
url: "https://routerai.ru/api/v1"
|
||||
first_byte_timeout: PT30S # окно для всех моделей провайдера
|
||||
|
||||
upstreams:
|
||||
- id: routerai-gpt4o
|
||||
provider: routerai
|
||||
model: gpt-4o
|
||||
first_byte_timeout: PT0S # у модели watchdog отключён
|
||||
```
|
||||
|
||||
При старте выводится, что настроено:
|
||||
`[llm-proxy] first_byte_timeout: upstreams=routerai-gpt4o=PT0S providers=routerai=PT30S default=PT30S`.
|
||||
|
||||
### Фоновый пинг апстримов в откате (`probe_interval`)
|
||||
|
||||
Когда апстрим «уходит в аут» надолго (backoff растёт до потолка), без пинга
|
||||
он вернётся в ротацию только по истечении интервала отката — это может быть
|
||||
минуты простоя, хотя модель уже могла ожить. `probe_interval` включает
|
||||
фоновый «пингер»: пока апстрим в backoff-откате, HealthProber каждые
|
||||
`probe_interval` шлёт ему минимальный запрос (`2+2=?`, `max_tokens=4`,
|
||||
`stream=true`) и слушает первый байт ответа. Как только апстрим ответил —
|
||||
его снимают с отката немедленно (`backoff.recordSuccess`), и он возвращается
|
||||
в общий пул, не дожидаясь истечения backoff-интервала.
|
||||
|
||||
| Поле | Тип / дефолт | Значение |
|
||||
|---|---|---|
|
||||
| `providers[].probe_interval` | ISO-8601-длительность / отсутствует | Интервал пинга для всех моделей провайдера |
|
||||
| `upstreams[].probe_interval` | ISO-8601-длительность / отсутствует | Интервал для конкретной модели (приоритет над провайдерским) |
|
||||
|
||||
**Дефолт** (если не задан нигде): `PT1S`.
|
||||
|
||||
**Приоритет — модели.** Если `upstreams[].probe_interval` задан, он
|
||||
перекрывает провайдерский; если нет — берётся провайдерский; если нет и там —
|
||||
дефолт `PT1S`.
|
||||
|
||||
**Отключение:** `PT0S` на провайдере или апстриме отключает пинг для этого
|
||||
апстрима — тогда он возвращается в пул только по истечении backoff-интервала.
|
||||
|
||||
Особенности:
|
||||
|
||||
- пингуются **только** апстримы, которые сейчас в backoff-откате; живой
|
||||
апстрим пинг не получает (проверка состояния — каждые `probe_interval`);
|
||||
- пинг **не занимает слот конкурентности** (идёт в обход `max_concurrency`),
|
||||
чтобы не отжимать слот у живых запросов;
|
||||
- ошибка пинга **не** считается в backoff (иначе каждая попытка удваивала
|
||||
бы интервал отката); учитывается только успех;
|
||||
- признак «ответил» — HTTP 2xx и первый байт тела за `first_byte_timeout`
|
||||
(тот же watchdog, что и у обычных запросов);
|
||||
- на пинг применяются патчи провайдера и апстрима (как к обычному запросу),
|
||||
чтобы он валидировал реальный путь;
|
||||
- при восстановлении — лог:
|
||||
`[llm-proxy] probe upstream=<id> ответил — снят с backoff, вернулся в пул`,
|
||||
метрика `llm_proxy_probe_pings_total{result="ok"}`;
|
||||
- при молчании — лог `[llm-proxy] probe upstream=<id> молчит — остаюсь в backoff`,
|
||||
метрика `llm_proxy_probe_pings_total{result="fail"}`.
|
||||
|
||||
Значение — ISO-8601-длительность, формат как у `backoff` (`PT1S`, `PT500MS`,
|
||||
`PT5S`, …).
|
||||
|
||||
```yaml
|
||||
providers:
|
||||
- id: routerai
|
||||
url: "https://routerai.ru/api/v1"
|
||||
probe_interval: PT1S # пингуем все модели провайдера раз в секунду
|
||||
|
||||
upstreams:
|
||||
- id: routerai-gpt4o
|
||||
provider: routerai
|
||||
model: gpt-4o
|
||||
probe_interval: PT0S # у модели пинг отключён (ждать истечения backoff)
|
||||
```
|
||||
|
||||
При старте выводится, что настроено:
|
||||
`[llm-proxy] probe_interval: upstreams=routerai-gpt4o=PT0S providers=routerai=PT1S default=PT1S`.
|
||||
|
||||
### Prometheus-метрики (`/metrics`)
|
||||
|
||||
Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics`
|
||||
@@ -328,7 +460,8 @@ upstreams:
|
||||
|
||||
| Метрика | Тип | Смысл |
|
||||
|---|---|---|
|
||||
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `cancelled` — клиент отвалился посреди стрима |
|
||||
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `timeout` — апстрим не прислал первый байт стрима за `first_byte_timeout`, `cancelled` — клиент отвалился посреди стрима |
|
||||
| `llm_proxy_probe_pings_total{upstream, provider, result}` | counter | фоновые пинги апстримов в backoff-откате (HealthProber). `result`: `ok` — апстрим ответил (снят с отката), `fail` — молчит (остался в откате) |
|
||||
| `llm_proxy_upstream_inflight{upstream, provider}` | gauge | занятые слоты апстрима прямо сейчас (конкурентность) |
|
||||
| `llm_proxy_upstream_fail_streak{upstream, provider}` | gauge | счётчик сбоев подряд (backoff): сколько раз подряд упал |
|
||||
| `llm_proxy_upstream_cooling_seconds{upstream, provider}` | gauge | сколько секунд апстрим ещё в backoff-откате (0 = жив) |
|
||||
@@ -340,7 +473,7 @@ upstreams:
|
||||
- оборот по виртуальным моделям: `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)`.
|
||||
- доля ошибок провайдера: `sum(rate(llm_proxy_requests_total{result=~"4xx|429|402|5xx|net_err|timeout"}[10m])) by (provider) / sum(rate(llm_proxy_requests_total[10m])) by (provider)`.
|
||||
|
||||
### Пример сборки тела (многослойный `patch`)
|
||||
|
||||
@@ -426,10 +559,11 @@ data class ProviderConf(
|
||||
val think_tags: String? = null,
|
||||
val reasoning_field: String? = null,
|
||||
val reasoning_empty_ok: Boolean = false,
|
||||
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
|
||||
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
|
||||
val first_byte_timeout: Duration? = null, // окно первого байта стрима (ISO-8601)
|
||||
val probe_interval: Duration? = null, // интервал фонового пинга (ISO-8601)
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class UpstreamConf(
|
||||
val id: String, // наш внутренний id апстрима
|
||||
val provider: String, // ссылка на ProviderConf.id
|
||||
@@ -438,6 +572,8 @@ data class UpstreamConf(
|
||||
val patch: JsonObject? = null, // к запросам этой апстрим-модели
|
||||
val think_tags: String? = null, // переопределение think-режима модели
|
||||
val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским)
|
||||
val first_byte_timeout: Duration? = null, // окно первого байта стрима (приоритет над провайдерским)
|
||||
val probe_interval: Duration? = null, // интервал фонового пинга (приоритет над провайдерским)
|
||||
)
|
||||
|
||||
@Serializable
|
||||
@@ -602,6 +738,10 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
|
||||
(`(backoff: отдых PT…)`); при пропуске апстрима в откате —
|
||||
`upstream=<id> в backoff-откате (~Ns) — пропускаю`; при старте — список
|
||||
настроенных backoff (`backoff: upstreams=… providers=…`).
|
||||
- **Таймаут первого байта стрима**: при срабатывании watchdog'а —
|
||||
`upstream=<id> TIMEOUT: первый байт не пришёл за <N>ms → фейловер`;
|
||||
при старте — список настроенных `first_byte_timeout`
|
||||
(`first_byte_timeout: upstreams=… providers=… default=PT30S`).
|
||||
- **Все в откате** — `503 all upstreams cooling, retry in Ns` для
|
||||
`model=<витрина>` с заголовком `Retry-After: N`.
|
||||
|
||||
|
||||
@@ -30,6 +30,20 @@ non-stream.
|
||||
1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате
|
||||
— `503` с `Retry-After`.
|
||||
|
||||
Для `stream=true` запросов действует watchdog на первый байт ответа:
|
||||
`first_byte_timeout` (ISO-8601, дефолт `PT30S`, `PT0S` отключает) — если
|
||||
апстрим не прислал ни одного байта за окно, слот освобождается, ошибка идёт
|
||||
в backoff и роутер фейловерит на следующего свободного апстрима. Приоритет:
|
||||
`upstreams[].first_byte_timeout` → `providers[].first_byte_timeout` → `PT30S`.
|
||||
|
||||
Пока апстрим «отдыхает» в backoff-откате, его можно фоново пинговать:
|
||||
`probe_interval` (ISO-8601, дефолт `PT1S`, `PT0S` отключает) — HealthProber
|
||||
шлёт апстриму минимальный запрос (`2+2=?`, `max_tokens=4`) и, как только тот
|
||||
ответил, снимает его с отката немедленно, не дожидаясь истечения
|
||||
backoff-интервала. Пинг не занимает слот конкурентности и не засчитывается
|
||||
в backoff. Приоритет: `upstreams[].probe_interval` →
|
||||
`providers[].probe_interval` → `PT1S`.
|
||||
|
||||
Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по
|
||||
`model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) —
|
||||
подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список
|
||||
|
||||
+15
-1
@@ -14,7 +14,7 @@
|
||||
| --- | --- |
|
||||
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
|
||||
| `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
|
||||
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, неприкосновенность сообщений без `tool_calls` (в т.ч. `[]`) и `user`/`tool`, пустая строка при `emptyOk`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
|
||||
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, заполнение поля на **каждом** assistant-сообщении (с `tool_calls`, без, `tool_calls: []` — без текста и `emptyOk=false` не трогаем; `emptyOk=true` — пустая строка), неприкосновенность `user`/`tool`/`system`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
|
||||
| `NullToleranceTest` | Устойчивость разбора ответов апстрима к JSON-`null` (`choices`, `delta.tool_calls`, `delta.content`) — пропуск вместо исключения; |
|
||||
| `ThinkTagTransformTest` | Non-stream путь `transformThinkMessage`: перенос рассуждений в `reasoning_content`, дописывание к уже имеющемуся, strip, незакрытый блок, отсутствие изменений → null; |
|
||||
| `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; |
|
||||
@@ -83,6 +83,20 @@ Console Go; тело — как у opencode: assistant + `reasoning` + `reasonin
|
||||
`assistantWithoutToolCallsIsUntouched` и `assistantWithEmptyToolCallsIsUntouched`;
|
||||
возврат `obj["choices"]?.jsonArray` в `rebuildFromChunks` → падает `rebuildFromChunksToleratesNullChoices`.
|
||||
|
||||
**Поправка после боевого 400 (2026-09-13):** живые stateful-сессии Console Go
|
||||
(сессия, где глюч уже видел thinking-ответы) всё-таки требовали
|
||||
`reasoning_content` и на assistant-сообщениях **без** `tool_calls` — свежая
|
||||
сессия такой истории прощала (все формы выше — 200). Опция «на каждом
|
||||
assistant-сообщении» совпадает с тем, что делает сам opencode («Deepseek
|
||||
requires all assistant messages to have reasoning on them», пустая
|
||||
reasoning-часть дописывается тоже). Охрана `tool_calls` в
|
||||
`applyReasoningField` снята: поле дописывается на каждом assistant-сообщении
|
||||
(текст — из `reasoning`/`reasoning_details`, при `emptyOk` — пустая строка).
|
||||
Тесты `assistantWithoutToolCallsIsUntouched`/`assistantWithEmptyToolCallsIsUntouched`
|
||||
заменены на `assistantWithoutToolCallsAlsoGetsField`,
|
||||
`assistantWithoutReasoningAndNoEmptyOkIsUntouched`,
|
||||
`assistantWithoutToolCallsWithEmptyOkFillsEmpty`, `assistantWithEmptyToolCallsGetsField`.
|
||||
|
||||
## Известное ограничение
|
||||
|
||||
Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.call.body
|
||||
import io.ktor.client.request.headers
|
||||
import io.ktor.client.request.preparePost
|
||||
import io.ktor.client.request.setBody
|
||||
import io.ktor.utils.io.ByteReadChannel
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.currentCoroutineContext
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.serialization.json.JsonArray
|
||||
import kotlinx.serialization.json.JsonObject
|
||||
import kotlinx.serialization.json.JsonPrimitive
|
||||
import kotlin.time.Duration
|
||||
|
||||
/**
|
||||
* Фоновый «пингер» апстримов, которые в данный момент в backoff-откате.
|
||||
*
|
||||
* Пока апстрим молчит (в откате), HealthProber раз в [interval] шлёт ему
|
||||
* минимальный chat.completion (`2+2=?`, `max_tokens=4`) и слушает первый байт
|
||||
* ответа. Как только апстрим ответил — вызывается [BackoffRegistry.recordSuccess],
|
||||
* и он возвращается в общий пул сразу, не дожидаясь истечения backoff-интервала.
|
||||
* Это ускоряет восстановление после долгого простоя провайдера.
|
||||
*
|
||||
* Пинг НЕ занимает слот конкурентности ([UpstreamCounter]) — он идёт в обход,
|
||||
* чтобы не отжимать слот у живых запросов. Ошибки пинга НЕ считаются в backoff
|
||||
* (иначе каждая попытка удваивала бы интервал отката); считается только успех.
|
||||
*/
|
||||
class HealthProber(
|
||||
private val providersById: Map<String, ProviderConf>,
|
||||
private val upstreams: List<UpstreamConf>,
|
||||
private val http: HttpClient,
|
||||
private val backoff: BackoffRegistry,
|
||||
private val metrics: MetricsRegistry,
|
||||
private val scope: CoroutineScope,
|
||||
) {
|
||||
private val jobs = mutableMapOf<String, Job>()
|
||||
|
||||
/** Запустить пинг-корутины для всех апстримов с включённым пингом. */
|
||||
fun start() {
|
||||
for (up in upstreams) {
|
||||
val interval = effectiveProbeInterval(up, providersById[up.provider])
|
||||
if (interval != null) {
|
||||
jobs[up.id] = scope.launch { probeLoop(up, interval) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Цикл пинга одного апстрима: пингуем только пока он в откате. */
|
||||
internal suspend fun probeLoop(up: UpstreamConf, interval: Duration, pingImpl: suspend (UpstreamConf) -> Boolean = ::ping) {
|
||||
while (currentCoroutineContext().isActive) {
|
||||
if (backoff.isCooling(up)) {
|
||||
val healthy = pingImpl(up)
|
||||
metrics.recordProbe(up, if (healthy) "ok" else "fail")
|
||||
if (healthy) {
|
||||
backoff.recordSuccess(up)
|
||||
log.info { "[llm-proxy] probe upstream=${up.id} ответил — снят с backoff, вернулся в пул" }
|
||||
} else {
|
||||
log.info { "[llm-proxy] probe upstream=${up.id} молчит — остаюсь в backoff" }
|
||||
}
|
||||
}
|
||||
delay(interval)
|
||||
}
|
||||
}
|
||||
|
||||
/** Одна пинг-попытка: минимальный stream-запрос, слушаем первый байт. */
|
||||
private suspend fun ping(up: UpstreamConf): Boolean {
|
||||
val provider = providersById[up.provider] ?: return false
|
||||
val url = provider.url.trimEnd('/') + "/chat/completions"
|
||||
val providerKey = resolveEnv(provider.key)
|
||||
val body = buildProbeBody(up, provider)
|
||||
val firstByteTimeout = effectiveFirstByteTimeout(up, provider) ?: DEFAULT_FIRST_BYTE_TIMEOUT
|
||||
return try {
|
||||
http.preparePost(url) {
|
||||
headers {
|
||||
if (providerKey.isNotEmpty()) {
|
||||
set("Authorization", "Bearer $providerKey")
|
||||
}
|
||||
set("Content-Type", "application/json")
|
||||
}
|
||||
setBody(body.toString())
|
||||
}.execute { resp ->
|
||||
if (resp.status.value !in 200..299) {
|
||||
return@execute false
|
||||
}
|
||||
val ch = resp.body<ByteReadChannel>()
|
||||
val buf = ByteArray(64)
|
||||
readAvailableWithTimeout(ch, buf, firstByteTimeout.inWholeMilliseconds) > 0
|
||||
}
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
log.info { "[llm-proxy] probe upstream=${up.id} ошибка: ${e.message}" }
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Собрать тело пинг-запроса: минимальный chat.completion, на который любой
|
||||
* живой апстрим ответит сразу. Применяются патчи провайдера и апстрима (в том
|
||||
* же порядке, что и для реальных запросов), чтобы пинг валидировал реальный
|
||||
* путь запроса, но без клиентских данных и модели.
|
||||
*/
|
||||
internal fun buildProbeBody(up: UpstreamConf, provider: ProviderConf?): JsonObject {
|
||||
val base = JsonObject(
|
||||
mapOf(
|
||||
"model" to JsonPrimitive(up.model),
|
||||
"messages" to JsonArray(
|
||||
listOf(
|
||||
JsonObject(
|
||||
mapOf(
|
||||
"role" to JsonPrimitive("user"),
|
||||
"content" to JsonPrimitive("2+2=?"),
|
||||
),
|
||||
),
|
||||
),
|
||||
),
|
||||
"max_tokens" to JsonPrimitive(4),
|
||||
"stream" to JsonPrimitive(true),
|
||||
),
|
||||
)
|
||||
var acc = base
|
||||
listOf(provider?.patch, up.patch).filterNotNull().forEach { patch ->
|
||||
acc = merge(acc, patch)
|
||||
}
|
||||
return acc
|
||||
}
|
||||
@@ -28,6 +28,14 @@ import io.ktor.utils.io.ByteReadChannel
|
||||
import io.ktor.utils.io.ByteWriteChannel
|
||||
import io.ktor.utils.io.writeFully
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.cancelAndJoin
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.selects.onTimeout
|
||||
import kotlinx.coroutines.selects.select
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.datetime.Clock
|
||||
import kotlinx.serialization.Serializable
|
||||
@@ -103,6 +111,20 @@ fun main() {
|
||||
" providers=" +
|
||||
config.providers.filter { it.backoff != null }.map { p -> "${p.id}=${p.backoff?.toIsoString()}" }.joinToString(", ")
|
||||
}
|
||||
log.info {
|
||||
"[llm-proxy] first_byte_timeout: upstreams=" +
|
||||
config.upstreams.filter { it.first_byte_timeout != null }.map { up -> "${up.id}=${up.first_byte_timeout?.toIsoString()}" }.joinToString(", ") +
|
||||
" providers=" +
|
||||
config.providers.filter { it.first_byte_timeout != null }.map { p -> "${p.id}=${p.first_byte_timeout?.toIsoString()}" }.joinToString(", ") +
|
||||
" default=${DEFAULT_FIRST_BYTE_TIMEOUT.toIsoString()}"
|
||||
}
|
||||
log.info {
|
||||
"[llm-proxy] probe_interval: upstreams=" +
|
||||
config.upstreams.filter { it.probe_interval != null }.map { up -> "${up.id}=${up.probe_interval?.toIsoString()}" }.joinToString(", ") +
|
||||
" providers=" +
|
||||
config.providers.filter { it.probe_interval != null }.map { p -> "${p.id}=${p.probe_interval?.toIsoString()}" }.joinToString(", ") +
|
||||
" default=${DEFAULT_PROBE_INTERVAL.toIsoString()}"
|
||||
}
|
||||
log.info {
|
||||
"[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}"
|
||||
}
|
||||
@@ -112,6 +134,8 @@ fun main() {
|
||||
|
||||
val http = createHttpClient()
|
||||
val sessions = SessionRegistry()
|
||||
val prober = HealthProber(providersById, config.upstreams, http, backoff, metrics, CoroutineScope(SupervisorJob() + Dispatchers.Default))
|
||||
prober.start()
|
||||
startServer(config.server.host, config.server.port) {
|
||||
proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
|
||||
}
|
||||
@@ -279,13 +303,14 @@ private suspend fun handleChat(
|
||||
val ct = resp.headers["Content-Type"] ?: "text/event-stream"
|
||||
val status = HttpStatusCode.fromValue(upstreamStatus)
|
||||
val thinkMode = effectiveThinkTags(up, provider)
|
||||
val firstByteTimeout = effectiveFirstByteTimeout(up, provider)
|
||||
call.respondBytesWriter(ContentType.parse(ct), status) {
|
||||
val ch = resp.body<ByteReadChannel>()
|
||||
if (thinkMode == "off") {
|
||||
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
|
||||
streamRawWithDoneContract(ch)
|
||||
streamRawWithDoneContract(ch, firstByteTimeout)
|
||||
} else {
|
||||
streamSseWithThinkTags(ch, thinkMode)
|
||||
streamSseWithThinkTags(ch, thinkMode, firstByteTimeout)
|
||||
}
|
||||
}
|
||||
metrics.record(modelName, up, "ok")
|
||||
@@ -322,6 +347,15 @@ private suspend fun handleChat(
|
||||
|
||||
if (failover) continue
|
||||
if (responded) return
|
||||
} catch (e: FirstByteTimeoutException) {
|
||||
val rest = backoff.recordFailure(up)
|
||||
metrics.record(modelName, up, "timeout")
|
||||
log.warn {
|
||||
"[llm-proxy] chat model=$modelName upstream=${up.id} TIMEOUT: первый байт не пришёл за ${e.timeoutMs}ms → фейловер" +
|
||||
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
|
||||
}
|
||||
failed.add(up.id)
|
||||
continue
|
||||
} catch (e: CancellationException) {
|
||||
metrics.record(modelName, up, "cancelled")
|
||||
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
|
||||
@@ -493,6 +527,38 @@ internal fun effectiveThinkTags(up: UpstreamConf, provider: ProviderConf?): Stri
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Эффективный first-byte-timeout для стрим-ответа: значение у апстрима, если
|
||||
* задано; иначе у провайдера; иначе [DEFAULT_FIRST_BYTE_TIMEOUT]. Возвращает
|
||||
* null, если явно отключено (`Duration.ZERO` — `PT0S`): тогда стрим
|
||||
* ждёт первого байта без лимита (в пределах общего
|
||||
* [UPSTREAM_REQUEST_TIMEOUT_MS]).
|
||||
*/
|
||||
internal fun effectiveFirstByteTimeout(up: UpstreamConf, provider: ProviderConf?): Duration? {
|
||||
val raw = up.first_byte_timeout ?: provider?.first_byte_timeout ?: DEFAULT_FIRST_BYTE_TIMEOUT
|
||||
return if (raw <= Duration.ZERO) null else raw
|
||||
}
|
||||
|
||||
/**
|
||||
* Эффективный интервал фонового пинга апстрима в backoff-откате: значение у
|
||||
* апстрима, если задано; иначе у провайдера; иначе [DEFAULT_PROBE_INTERVAL].
|
||||
* Возвращает null, если явно отключено (`Duration.ZERO` — `PT0S`): тогда
|
||||
* HealthProber этот апстрим не пингует (ждём естественного истечения бэкоффа).
|
||||
*/
|
||||
internal fun effectiveProbeInterval(up: UpstreamConf, provider: ProviderConf?): Duration? {
|
||||
val raw = up.probe_interval ?: provider?.probe_interval ?: DEFAULT_PROBE_INTERVAL
|
||||
return if (raw <= Duration.ZERO) null else raw
|
||||
}
|
||||
|
||||
/**
|
||||
* Сигнализирует, что апстрим не прислал ни одного байта тела стрим-ответа за
|
||||
* [timeoutMs] мс. Обрабатывается как обычный сбой (фейловер, backoff,
|
||||
* метрика `timeout`), но отделён от [CancellationException] (отмена клиентом)
|
||||
* и от generic `Exception` (сетевые ошибки).
|
||||
*/
|
||||
class FirstByteTimeoutException(val timeoutMs: Long) :
|
||||
RuntimeException("upstream did not send first byte within ${timeoutMs}ms")
|
||||
|
||||
/** Освободить слот апстрима (в finally по завершении проксирования). */
|
||||
internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) {
|
||||
active.getValue(up.id).release()
|
||||
@@ -710,11 +776,14 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
|
||||
/**
|
||||
* Достроить нативное поле рассуждений для апстримов, которые его требуют
|
||||
* (Console Go / deepseek в thinking-режиме): если у провайдера объявлено
|
||||
* `reasoningField`, то в каждом assistant-сообщении с непустым `tool_calls`
|
||||
* добавляем это поле, ЕСЛИ его там ещё нет. Текст берём из `reasoning`
|
||||
* `reasoningField`, то в каждом assistant-сообщении добавляем это поле,
|
||||
* ЕСЛИ его там ещё нет (не только в тех, что с tool_calls — deepseek
|
||||
* требует `reasoning_content` на КАЖДОМ ассистент-сообщении в thinking-режиме,
|
||||
* см. «must be passed back to the API»). Текст берём из `reasoning`
|
||||
* (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`).
|
||||
* Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста
|
||||
* пишем пустую строку, только если `emptyOk`.
|
||||
* пишем пустую строку, только если `emptyOk` — апстрим требует самого
|
||||
* НАЛИЧИЯ поля, даже пустого (так делает и сам opencode).
|
||||
* Тело возвращается без изменений (тот же объект), если менять нечего.
|
||||
*/
|
||||
internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Boolean): JsonObject {
|
||||
@@ -725,8 +794,6 @@ internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Bool
|
||||
val msg = el as? JsonObject ?: return@map el
|
||||
val role = (msg["role"] as? JsonPrimitive)?.takeIf { it.isString }?.content
|
||||
if (role != "assistant") return@map el
|
||||
val toolCalls = msg["tool_calls"] as? JsonArray
|
||||
if (toolCalls == null || toolCalls.isEmpty()) return@map el
|
||||
val existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content
|
||||
if (existing != null && existing.isNotEmpty()) return@map el
|
||||
val text = reasoningTextOf(msg)
|
||||
@@ -826,12 +893,26 @@ internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDO
|
||||
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
|
||||
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
|
||||
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
|
||||
* Если задан [firstByteTimeout] (non-null), первый байт тела должен прийти
|
||||
* за это время, иначе бросаем [FirstByteTimeoutException]; после первого
|
||||
* байта watchdog выключается (ждём ответа вплоть до общего
|
||||
* [UPSTREAM_REQUEST_TIMEOUT_MS]).
|
||||
*/
|
||||
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
|
||||
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(
|
||||
source: ByteReadChannel,
|
||||
firstByteTimeout: Duration? = null,
|
||||
) {
|
||||
val detector = StreamEndDetector()
|
||||
val buf = ByteArray(8192)
|
||||
val timeoutMs = firstByteTimeout?.inWholeMilliseconds ?: 0L
|
||||
var isFirstRead = timeoutMs > 0L
|
||||
while (true) {
|
||||
val n = source.readAvailable(buf)
|
||||
val n = if (isFirstRead) {
|
||||
isFirstRead = false
|
||||
readAvailableWithTimeout(source, buf, timeoutMs)
|
||||
} else {
|
||||
source.readAvailable(buf)
|
||||
}
|
||||
if (n == -1) break
|
||||
if (n > 0) {
|
||||
writeFully(buf, 0, n)
|
||||
@@ -852,14 +933,30 @@ internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteRead
|
||||
* записывается как `data: <json>\n\n` (событие-граница SSE), а при ошибке парса —
|
||||
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
|
||||
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
|
||||
* Если задан [firstByteTimeout] (non-null), первая строка должна прийти за это
|
||||
* время, иначе [FirstByteTimeoutException] (watchdog только на первой строке).
|
||||
*/
|
||||
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
|
||||
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(
|
||||
source: ByteReadChannel,
|
||||
thinkMode: String,
|
||||
firstByteTimeout: Duration? = null,
|
||||
) {
|
||||
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
|
||||
val addReasoning = thinkMode == "split"
|
||||
var sawDone = false
|
||||
var sawFinishReason = false
|
||||
val timeoutMs = firstByteTimeout?.inWholeMilliseconds ?: 0L
|
||||
var isFirstRead = timeoutMs > 0L
|
||||
|
||||
suspend fun readNext(): String? = if (isFirstRead) {
|
||||
isFirstRead = false
|
||||
readLineWithTimeout(source, timeoutMs)
|
||||
} else {
|
||||
source.readLine(LineEnding.Lenient)
|
||||
}
|
||||
|
||||
while (true) {
|
||||
val line = source.readLine(LineEnding.Lenient) ?: break
|
||||
val line = readNext() ?: break
|
||||
when {
|
||||
line.startsWith("data:") -> {
|
||||
val payload = line.removePrefix("data:").trim()
|
||||
@@ -912,6 +1009,49 @@ internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadCha
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Прочитать из [source] в [buf] с жёстким лимитом [timeoutMs]: если за это
|
||||
* время ни одного байта не пришло, бросить [FirstByteTimeoutException].
|
||||
* `select` отменяет «проигравшую» ветку; при таймауте явно отменяем и
|
||||
* дожидаемся чтение (`cancelAndJoin`), чтобы соединение с апстримом
|
||||
* корректно закрылось. Сентинел [Int.MIN_VALUE] невозможен как результат
|
||||
* [ByteReadChannel.readAvailable] (это `-1` для EOF или `>= 0` для байт).
|
||||
*/
|
||||
internal suspend fun readAvailableWithTimeout(source: ByteReadChannel, buf: ByteArray, timeoutMs: Long): Int =
|
||||
coroutineScope {
|
||||
val readJob = async { source.readAvailable(buf) }
|
||||
val result = select<Int> {
|
||||
readJob.onAwait { it }
|
||||
onTimeout(timeoutMs) { Int.MIN_VALUE }
|
||||
}
|
||||
if (result == Int.MIN_VALUE) {
|
||||
readJob.cancelAndJoin()
|
||||
throw FirstByteTimeoutException(timeoutMs)
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
/**
|
||||
* Прочитать строку из [source] с лимитом [timeoutMs]: если за это время ни
|
||||
* одной строки/байта не пришло, бросить [FirstByteTimeoutException].
|
||||
* [ByteReadChannel.readLine] возвращает `null` и при EOF, и (через `onTimeout`)
|
||||
* при таймауте; различаем по состоянию [readJob]: если ветка чтения не
|
||||
* завершилась к моменту возврата из `select` — таймаут (select отменил её).
|
||||
*/
|
||||
internal suspend fun readLineWithTimeout(source: ByteReadChannel, timeoutMs: Long): String? =
|
||||
coroutineScope {
|
||||
val readJob = async { source.readLine(LineEnding.Lenient) }
|
||||
val result = select<String?> {
|
||||
readJob.onAwait { it }
|
||||
onTimeout(timeoutMs) { null }
|
||||
}
|
||||
if (result == null && !readJob.isCompleted) {
|
||||
readJob.cancelAndJoin()
|
||||
throw FirstByteTimeoutException(timeoutMs)
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
/**
|
||||
* Пересобрать SSE-чанк: по каждому choice (ключ `index`, дефолт 0) взять
|
||||
* `delta.content` (только если это JSON-строка; массив частей не трогаем) и
|
||||
@@ -965,6 +1105,8 @@ data class ProviderConf(
|
||||
val reasoning_field: String? = null,
|
||||
val reasoning_empty_ok: Boolean = false,
|
||||
val backoff: Duration? = null,
|
||||
val first_byte_timeout: Duration? = null,
|
||||
val probe_interval: Duration? = null,
|
||||
)
|
||||
|
||||
data class UpstreamConf(
|
||||
@@ -975,6 +1117,8 @@ data class UpstreamConf(
|
||||
val patch: JsonObject? = null,
|
||||
val think_tags: String? = null,
|
||||
val backoff: Duration? = null,
|
||||
val first_byte_timeout: Duration? = null,
|
||||
val probe_interval: Duration? = null,
|
||||
)
|
||||
|
||||
data class ModelConf(
|
||||
@@ -1033,6 +1177,8 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
reasoning_field = m.strOrNull("reasoning_field"),
|
||||
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
|
||||
backoff = m.durationOrNull("backoff"),
|
||||
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
|
||||
probe_interval = m.durationOrNull("probe_interval"),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -1046,6 +1192,8 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
|
||||
think_tags = m.strOrNull("think_tags"),
|
||||
backoff = m.durationOrNull("backoff"),
|
||||
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
|
||||
probe_interval = m.durationOrNull("probe_interval"),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -9,7 +9,9 @@ import kotlinx.coroutines.sync.withLock
|
||||
*
|
||||
* Метрики (всё в памяти, сброс при рестарте):
|
||||
* - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик
|
||||
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | cancelled`;
|
||||
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | timeout | cancelled`;
|
||||
* - `llm_proxy_probe_pings_total{upstream,provider,result}` — счётчик фоновых
|
||||
* пингов апстримов в backoff-откате (HealthProber); `result`: `ok | fail`;
|
||||
* - `llm_proxy_upstream_inflight{upstream,provider}` — сейчас в работе (слоты);
|
||||
* - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев
|
||||
* (счётчик backoff);
|
||||
@@ -19,6 +21,7 @@ import kotlinx.coroutines.sync.withLock
|
||||
class MetricsRegistry {
|
||||
private val lock = Mutex()
|
||||
private val requests = HashMap<RequestKey, Long>()
|
||||
private val probes = HashMap<ProbeKey, Long>()
|
||||
|
||||
private data class RequestKey(
|
||||
val model: String,
|
||||
@@ -27,6 +30,12 @@ class MetricsRegistry {
|
||||
val result: String,
|
||||
)
|
||||
|
||||
private data class ProbeKey(
|
||||
val upstream: String,
|
||||
val provider: String,
|
||||
val result: String,
|
||||
)
|
||||
|
||||
/** Учёт chat-запроса по исходу (значения `result` — в описании класса). */
|
||||
suspend fun record(model: String, up: UpstreamConf, result: String) {
|
||||
lock.withLock {
|
||||
@@ -35,6 +44,14 @@ class MetricsRegistry {
|
||||
}
|
||||
}
|
||||
|
||||
/** Учёт пинга апстрима в backoff-откате (result: `ok`/`fail`). */
|
||||
suspend fun recordProbe(up: UpstreamConf, result: String) {
|
||||
lock.withLock {
|
||||
val key = ProbeKey(up.id, up.provider, result)
|
||||
probes[key] = (probes[key] ?: 0L) + 1
|
||||
}
|
||||
}
|
||||
|
||||
/** Текстовая выгрузка в формате Prometheus на текущий момент. */
|
||||
suspend fun render(
|
||||
config: Config,
|
||||
@@ -56,6 +73,19 @@ class MetricsRegistry {
|
||||
)
|
||||
}
|
||||
}
|
||||
sb.append("# HELP llm_proxy_probe_pings_total Фоновые пинги апстримов в backoff-откате\n")
|
||||
sb.append("# TYPE llm_proxy_probe_pings_total counter\n")
|
||||
lock.withLock {
|
||||
probes.entries
|
||||
.sortedWith(compareBy({ it.key.upstream }, { it.key.provider }, { it.key.result }))
|
||||
.forEach { (k, n) ->
|
||||
sb.append(
|
||||
"llm_proxy_probe_pings_total{upstream=\"" + esc(k.upstream) +
|
||||
"\",provider=\"" + esc(k.provider) +
|
||||
"\",result=\"" + esc(k.result) + "\"} $n\n",
|
||||
)
|
||||
}
|
||||
}
|
||||
sb.append("# HELP llm_proxy_upstream_inflight Занято слотов апстрима сейчас\n")
|
||||
sb.append("# TYPE llm_proxy_upstream_inflight gauge\n")
|
||||
for (up in config.upstreams) {
|
||||
|
||||
@@ -1,4 +1,26 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
/** Предел времени на весь запрос к апстриму, мс: от отправки до приёма всего ответа (включая стриминг). */
|
||||
const val UPSTREAM_REQUEST_TIMEOUT_MS: Long = 5 * 60_000L
|
||||
|
||||
/**
|
||||
* Default «watchdog» на первый байт тела стрим-ответа: если апстрим не прислал
|
||||
* ни одного байта за это время, прокси считает это сбоем апстрима и
|
||||
* фейловерит (освобождает слот конкурентности, считает ошибку в backoff и
|
||||
* идёт к следующему апстриму модели). Применяется ТОЛЬКО к stream=true
|
||||
* запросам; для non-stream работает общий [UPSTREAM_REQUEST_TIMEOUT_MS].
|
||||
* Отключается явным `first_byte_timeout: PT0S` на апстриме или провайдере.
|
||||
*/
|
||||
val DEFAULT_FIRST_BYTE_TIMEOUT: Duration = 30.seconds
|
||||
|
||||
/**
|
||||
* Default интервал фоновой проверки («пинга») апстрима, который сейчас в
|
||||
* backoff-откате: пока апстрим молчит, HealthProber шлёт ему минимальный
|
||||
* запрос (2+2=?, max_tokens=4) каждые [DEFAULT_PROBE_INTERVAL], и как только
|
||||
* он отвечает — снимает его с отката раньше, чем истечёт интервал бэкоффа.
|
||||
* Отключается явным `probe_interval: PT0S` на апстриме или провайдере.
|
||||
*/
|
||||
val DEFAULT_PROBE_INTERVAL: Duration = 1.seconds
|
||||
|
||||
@@ -4,6 +4,7 @@ import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
import io.ktor.http.headersOf
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.serialization.json.Json
|
||||
@@ -586,4 +587,60 @@ class ConfigLogicTest {
|
||||
assertEquals(null, cfg.providers[0].reasoning_field)
|
||||
assertEquals(false, cfg.providers[0].reasoning_empty_ok)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun parseConfigReadsFirstByteTimeoutOnProviderAndUpstream() {
|
||||
val yaml = """
|
||||
providers:
|
||||
- id: p1
|
||||
url: "https://x.ru/api/v1"
|
||||
first_byte_timeout: PT45S
|
||||
- id: p2
|
||||
url: "https://y.ru/api/v1"
|
||||
upstreams:
|
||||
- id: u1
|
||||
provider: p1
|
||||
model: real-1
|
||||
first_byte_timeout: PT0S
|
||||
- id: u2
|
||||
provider: p2
|
||||
model: real-2
|
||||
models:
|
||||
- name: m1
|
||||
upstreams: [u1, u2]
|
||||
""".trimIndent()
|
||||
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
|
||||
assertEquals(45.seconds, cfg.providers[0].first_byte_timeout)
|
||||
assertEquals(null, cfg.providers[1].first_byte_timeout)
|
||||
assertEquals(0.seconds, cfg.upstreams[0].first_byte_timeout)
|
||||
assertEquals(null, cfg.upstreams[1].first_byte_timeout)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun parseConfigReadsProbeIntervalOnProviderAndUpstream() {
|
||||
val yaml = """
|
||||
providers:
|
||||
- id: p1
|
||||
url: "https://x.ru/api/v1"
|
||||
probe_interval: PT2S
|
||||
- id: p2
|
||||
url: "https://y.ru/api/v1"
|
||||
upstreams:
|
||||
- id: u1
|
||||
provider: p1
|
||||
model: real-1
|
||||
probe_interval: PT0S
|
||||
- id: u2
|
||||
provider: p2
|
||||
model: real-2
|
||||
models:
|
||||
- name: m1
|
||||
upstreams: [u1, u2]
|
||||
""".trimIndent()
|
||||
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
|
||||
assertEquals(2.seconds, cfg.providers[0].probe_interval)
|
||||
assertEquals(null, cfg.providers[1].probe_interval)
|
||||
assertEquals(0.seconds, cfg.upstreams[0].probe_interval)
|
||||
assertEquals(null, cfg.upstreams[1].probe_interval)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,218 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import io.ktor.utils.io.ByteChannel
|
||||
import io.ktor.utils.io.close
|
||||
import io.ktor.utils.io.readAvailable
|
||||
import io.ktor.utils.io.writeFully
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.advanceTimeBy
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
|
||||
/**
|
||||
* Тесты watchdog'а на первый байт тела стрим-ответа.
|
||||
*
|
||||
* Сценарии:
|
||||
* - Апстрим не прислал ничего за N мс → [FirstByteTimeoutException] на первой
|
||||
* попытке чтения (после первой попытки watchdog отключается).
|
||||
* - Апстрим закрыл канал до таймаута, не прислав данных → возврат EOF (-1 /
|
||||
* null), НЕ [FirstByteTimeoutException].
|
||||
* - Апстрим прислал первый байт вовремя → читаем его, watchdog отключён.
|
||||
* - Когда таймаут не задан (`null`) → читаем без лимита (как и было).
|
||||
*/
|
||||
class FirstByteTimeoutTest {
|
||||
|
||||
/** Окно для watchdog'а в тестах: виртуальное время `runTest` его мгновенно прокручивает. */
|
||||
private val watchWindow = 100.milliseconds
|
||||
|
||||
@Test
|
||||
fun readAvailableTimesOutWhenNoBytesAndNotClosed() = runTest {
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
val buf = ByteArray(64)
|
||||
val done = CompletableDeferred<Throwable?>()
|
||||
|
||||
launch {
|
||||
try {
|
||||
readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
|
||||
done.complete(AssertionError("expected FirstByteTimeoutException"))
|
||||
} catch (e: FirstByteTimeoutException) {
|
||||
done.complete(null)
|
||||
assertEquals(watchWindow.inWholeMilliseconds, e.timeoutMs)
|
||||
} catch (e: Throwable) {
|
||||
done.complete(e)
|
||||
}
|
||||
}
|
||||
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
|
||||
runCurrent()
|
||||
val err = done.await()
|
||||
if (err != null) throw err
|
||||
}
|
||||
|
||||
@Test
|
||||
fun readAvailableReturnsEofImmediatelyWhenChannelClosedEmpty() = runTest {
|
||||
// Канал сразу закрыт без данных: -1 (EOF), никакого таймаута.
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
src.close(null)
|
||||
val buf = ByteArray(64)
|
||||
|
||||
val n = readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
|
||||
assertEquals(-1, n, "закрытый пустой канал → EOF, не таймаут")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun readAvailableReturnsBytesWhenDataArrivesInTime() = runTest {
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
val payload = "hello".encodeToByteArray()
|
||||
src.writeFully(payload)
|
||||
val buf = ByteArray(64)
|
||||
|
||||
val n = readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
|
||||
assertEquals(payload.size, n)
|
||||
assertEquals(payload.decodeToString(), buf.decodeToString(0, n))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun readLineTimesOutWhenNoLinesAndNotClosed() = runTest {
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
val done = CompletableDeferred<Throwable?>()
|
||||
|
||||
launch {
|
||||
try {
|
||||
readLineWithTimeout(src, watchWindow.inWholeMilliseconds)
|
||||
done.complete(AssertionError("expected FirstByteTimeoutException"))
|
||||
} catch (e: FirstByteTimeoutException) {
|
||||
done.complete(null)
|
||||
} catch (e: Throwable) {
|
||||
done.complete(e)
|
||||
}
|
||||
}
|
||||
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
|
||||
runCurrent()
|
||||
val err = done.await()
|
||||
if (err != null) throw err
|
||||
}
|
||||
|
||||
@Test
|
||||
fun readLineReturnsNullWhenChannelClosedEmpty() = runTest {
|
||||
// Закрытый пустой канал → null (EOF), не таймаут. Это критично: иначе бы
|
||||
// мы фейловерили легитимные короткие ответы, где апстрим закрыл стрим
|
||||
// до первой строки.
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
src.close(null)
|
||||
|
||||
val line = readLineWithTimeout(src, watchWindow.inWholeMilliseconds)
|
||||
assertEquals(null, line)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun rawStreamTimesOutWhenUpstreamSilent() = runTest {
|
||||
// Сырой passthrough: апстрим молчит вечно → на первой попытке чтения
|
||||
// вылетает FirstByteTimeoutException, в out ничего не пишется.
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
val out = ByteChannel(autoFlush = true)
|
||||
val done = CompletableDeferred<Throwable?>()
|
||||
|
||||
launch {
|
||||
try {
|
||||
out.streamRawWithDoneContract(src, watchWindow)
|
||||
done.complete(AssertionError("expected FirstByteTimeoutException"))
|
||||
} catch (e: FirstByteTimeoutException) {
|
||||
done.complete(null)
|
||||
} catch (e: Throwable) {
|
||||
done.complete(e)
|
||||
}
|
||||
}
|
||||
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
|
||||
runCurrent()
|
||||
val err = done.await()
|
||||
if (err != null) throw err
|
||||
|
||||
// Подтверждаем, что в out ничего не утекло.
|
||||
out.close(null)
|
||||
val buf = ByteArray(64)
|
||||
val n = out.readAvailable(buf)
|
||||
assertEquals(-1, n, "до таймаута клиенту ничего не должно уйти")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun thinkStreamTimesOutWhenUpstreamSilent() = runTest {
|
||||
// SSE-вариант: первая строка не приходит → FirstByteTimeoutException,
|
||||
// в out ничего не пишется.
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
val out = ByteChannel(autoFlush = true)
|
||||
val done = CompletableDeferred<Throwable?>()
|
||||
|
||||
launch {
|
||||
try {
|
||||
out.streamSseWithThinkTags(src, "split", watchWindow)
|
||||
done.complete(AssertionError("expected FirstByteTimeoutException"))
|
||||
} catch (e: FirstByteTimeoutException) {
|
||||
done.complete(null)
|
||||
} catch (e: Throwable) {
|
||||
done.complete(e)
|
||||
}
|
||||
}
|
||||
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
|
||||
runCurrent()
|
||||
val err = done.await()
|
||||
if (err != null) throw err
|
||||
|
||||
out.close(null)
|
||||
val buf = ByteArray(64)
|
||||
val n = out.readAvailable(buf)
|
||||
assertEquals(-1, n)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun thinkStreamFirstLineInTimeThenReadsUnlimited() = runTest {
|
||||
// Первый байт/строка пришли вовремя → watchdog выключен, остальной
|
||||
// поток читается обычным порядком. Проверяем, что НЕ бросило
|
||||
// FirstByteTimeoutException и данные дошли до out.
|
||||
val src = ByteChannel(autoFlush = true)
|
||||
val out = ByteChannel(autoFlush = true)
|
||||
val payload = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
|
||||
src.writeFully(payload.encodeToByteArray())
|
||||
src.close(null)
|
||||
|
||||
out.streamSseWithThinkTags(src, "split", watchWindow)
|
||||
|
||||
out.close(null)
|
||||
val sb = StringBuilder()
|
||||
val buf = ByteArray(256)
|
||||
while (true) {
|
||||
val n = out.readAvailable(buf)
|
||||
if (n == -1) break
|
||||
if (n > 0) sb.append(buf.decodeToString(0, n))
|
||||
}
|
||||
assertTrue(sb.toString().contains("data:"), "данные дошли до out: $sb")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun effectiveFirstByteTimeoutPrefersUpstreamThenProviderThenDefault() {
|
||||
// upstream > provider > default
|
||||
val provider = ProviderConf("p", "https://x", "", null, null, first_byte_timeout = 20.milliseconds)
|
||||
val upstream = UpstreamConf("u", "p", "m", null, null, null, null, first_byte_timeout = 5.milliseconds)
|
||||
val upstreamNoOwn = UpstreamConf("u", "p", "m")
|
||||
|
||||
assertEquals(5.milliseconds, effectiveFirstByteTimeout(upstream, provider))
|
||||
assertEquals(20.milliseconds, effectiveFirstByteTimeout(upstreamNoOwn, provider))
|
||||
assertEquals(DEFAULT_FIRST_BYTE_TIMEOUT, effectiveFirstByteTimeout(upstreamNoOwn, null))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun effectiveFirstByteTimeoutZeroDisablesWatchdog() {
|
||||
// 0s / PT0S → null: passthrough без watchdog'а.
|
||||
val up = UpstreamConf("u", "p", "m", null, null, null, null, first_byte_timeout = 0.milliseconds)
|
||||
assertEquals(null, effectiveFirstByteTimeout(up, null))
|
||||
|
||||
val providerZero = ProviderConf("p", "https://x", "", null, null, first_byte_timeout = 0.milliseconds)
|
||||
val upNoOwn = UpstreamConf("u", "p", "m")
|
||||
assertEquals(null, effectiveFirstByteTimeout(upNoOwn, providerZero))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.TestScope
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.json.jsonArray
|
||||
import kotlinx.serialization.json.jsonObject
|
||||
import kotlinx.serialization.json.jsonPrimitive
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.minutes
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
class HealthProberTest {
|
||||
|
||||
private val json = Json { ignoreUnknownKeys = true }
|
||||
|
||||
@Test
|
||||
fun effectiveProbeIntervalPrefersUpstreamThenProviderThenDefault() {
|
||||
val provider = ProviderConf("p", "https://x", "", null, null, probe_interval = 2.seconds)
|
||||
val upstream = UpstreamConf("u", "p", "m", null, null, null, null, null, probe_interval = 5.milliseconds)
|
||||
val upstreamNoOwn = UpstreamConf("u2", "p", "m")
|
||||
assertEquals(5.milliseconds, effectiveProbeInterval(upstream, provider))
|
||||
assertEquals(2.seconds, effectiveProbeInterval(upstreamNoOwn, provider))
|
||||
assertEquals(DEFAULT_PROBE_INTERVAL, effectiveProbeInterval(upstreamNoOwn, null))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun effectiveProbeIntervalZeroDisablesProber() {
|
||||
val up = UpstreamConf("u", "p", "m", null, null, null, null, null, probe_interval = 0.milliseconds)
|
||||
assertNull(effectiveProbeInterval(up, null))
|
||||
|
||||
val providerZero = ProviderConf("p", "https://x", "", null, null, probe_interval = 0.milliseconds)
|
||||
assertNull(effectiveProbeInterval(UpstreamConf("u2", "p", "m"), providerZero))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun buildProbeBodyIsMinimalChatWithModelReplacedAndPatchesApplied() {
|
||||
val provider = ProviderConf(
|
||||
"p", "https://x", "", null,
|
||||
Json.parseToJsonElement("""{"provider":{"allow_fallbacks":false}}""").jsonObject,
|
||||
probe_interval = 1.seconds,
|
||||
)
|
||||
val up = UpstreamConf(
|
||||
"u", "p", "real-1", null,
|
||||
Json.parseToJsonElement("""{"temperature":0.5}""").jsonObject,
|
||||
null, null, null, 1.seconds,
|
||||
)
|
||||
|
||||
val body = buildProbeBody(up, provider)
|
||||
|
||||
assertEquals("real-1", body["model"]?.jsonPrimitive?.content)
|
||||
assertEquals("2+2=?", body["messages"]?.jsonArray?.first()?.jsonObject?.get("content")?.jsonPrimitive?.content)
|
||||
assertEquals("true", body["stream"]?.jsonPrimitive?.content)
|
||||
assertEquals("4", body["max_tokens"]?.jsonPrimitive?.content)
|
||||
assertEquals("false", body["provider"]?.jsonObject?.get("allow_fallbacks")?.jsonPrimitive?.content)
|
||||
assertEquals("0.5", body["temperature"]?.jsonPrimitive?.content)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun buildProbeBodyIgnoresModelPatch() {
|
||||
val up = UpstreamConf("u", "p", "real-1")
|
||||
val provider = ProviderConf("p", "https://x")
|
||||
val body = buildProbeBody(up, provider)
|
||||
assertEquals(
|
||||
"""{"model":"real-1","messages":[{"role":"user","content":"2+2=?"}],"max_tokens":4,"stream":true}""",
|
||||
body.toString(),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun probeLoopRecoversUpstreamWhenPingSucceeds() = runTest {
|
||||
val provider = ProviderConf("p", "https://x", "", null, null, backoff = 5.minutes)
|
||||
val up = UpstreamConf("u", "p", "m", null, null, null, backoff = 5.minutes, probe_interval = 5.milliseconds)
|
||||
val backoff = BackoffRegistry(mapOf("p" to provider), mapOf("u" to up))
|
||||
backoff.recordFailure(up)
|
||||
assertEquals(true, backoff.isCooling(up))
|
||||
|
||||
val metrics = MetricsRegistry()
|
||||
val prober = HealthProber(
|
||||
providersById = mapOf("p" to provider),
|
||||
upstreams = listOf(up),
|
||||
http = createHttpClient(),
|
||||
backoff = backoff,
|
||||
metrics = metrics,
|
||||
scope = this,
|
||||
)
|
||||
var pings = 0
|
||||
val job = launch {
|
||||
prober.probeLoop(up, 5.milliseconds) { pinged ->
|
||||
pings++
|
||||
assertEquals("u", pinged.id)
|
||||
true
|
||||
}
|
||||
}
|
||||
runUntil { pings >= 1 }
|
||||
// успех уже зафиксирован (recordSuccess) — даём циклу несколько
|
||||
// итераций и убеждаемся, что повторных пингов не было
|
||||
delay(10.milliseconds)
|
||||
job.cancel()
|
||||
|
||||
assertEquals(1, pings)
|
||||
assertEquals(false, backoff.isCooling(up))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun probeLoopKeepsCoolingWhenPingFails() = runTest {
|
||||
val provider = ProviderConf("p", "https://x", "", null, null, backoff = 5.minutes)
|
||||
val up = UpstreamConf("u", "p", "m", null, null, null, backoff = 5.minutes, probe_interval = 5.milliseconds)
|
||||
val backoff = BackoffRegistry(mapOf("p" to provider), mapOf("u" to up))
|
||||
backoff.recordFailure(up)
|
||||
|
||||
val metrics = MetricsRegistry()
|
||||
val prober = HealthProber(
|
||||
providersById = mapOf("p" to provider),
|
||||
upstreams = listOf(up),
|
||||
http = createHttpClient(),
|
||||
backoff = backoff,
|
||||
metrics = metrics,
|
||||
scope = this,
|
||||
)
|
||||
var pings = 0
|
||||
val job = launch {
|
||||
prober.probeLoop(up, 5.milliseconds) { _ ->
|
||||
pings++
|
||||
false
|
||||
}
|
||||
}
|
||||
runUntil { pings >= 2 }
|
||||
job.cancel()
|
||||
|
||||
assertEquals(true, backoff.isCooling(up))
|
||||
}
|
||||
|
||||
private suspend fun TestScope.runUntil(condition: () -> Boolean) {
|
||||
var i = 0
|
||||
while (!condition() && i < 10_000) {
|
||||
delay(1)
|
||||
i++
|
||||
}
|
||||
if (!condition()) throw AssertionError("condition not met in time")
|
||||
}
|
||||
|
||||
}
|
||||
@@ -74,20 +74,42 @@ class ReasoningFieldTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun assistantWithoutToolCallsIsUntouched() {
|
||||
// Без tool_calls assistant-сообщение не трогаем (тело идентично).
|
||||
fun assistantWithoutToolCallsAlsoGetsField() {
|
||||
// DeepSeek требует reasoning_content на КАЖДОМ assistant-сообщении
|
||||
// (thinking-режим), а не только на тех, что с tool_calls.
|
||||
val m = """{"role":"assistant","reasoning":"R","content":"hi"}"""
|
||||
val out = apply(m, "reasoning_content", false)
|
||||
assertEquals("R", str(outMsgs(out)[0], "reasoning_content"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun assistantWithoutReasoningAndNoEmptyOkIsUntouched() {
|
||||
// assistant-сообщение без рассуждений и без tool_calls:
|
||||
// текста нет, emptyOk=false — поле не добавляем.
|
||||
val m = """{"role":"assistant","content":"hi"}"""
|
||||
val out = apply(m, "reasoning_content", false)
|
||||
assertEquals("""{"messages":[$m]}""", out.toString())
|
||||
assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun assistantWithEmptyToolCallsIsUntouched() {
|
||||
// Пустой массив tool_calls = не трогаем (тело идентично).
|
||||
fun assistantWithoutToolCallsWithEmptyOkFillsEmpty() {
|
||||
// Без tool_calls и без текста, но emptyOk=true — дописываем пустую строку,
|
||||
// т.к. апстриму важно само НАЛИЧИЕ поля.
|
||||
val m = """{"role":"assistant","content":"hi"}"""
|
||||
val out = apply(m, "reasoning_content", true)
|
||||
val m0 = outMsgs(out)[0]
|
||||
assertEquals("", str(m0, "reasoning_content"))
|
||||
assertTrue(m0.containsKey("reasoning_content"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun assistantWithEmptyToolCallsGetsField() {
|
||||
// Пустой массив tool_calls больше не исключает сообщение из обработки:
|
||||
// текст из reasoning дописывается в нативное поле.
|
||||
val m = """{"role":"assistant","tool_calls":[],"reasoning":"R"}"""
|
||||
val out = apply(m, "reasoning_content", false)
|
||||
assertEquals("""{"messages":[$m]}""", out.toString())
|
||||
assertEquals("R", str(outMsgs(out)[0], "reasoning_content"))
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user