2 Commits
9 .. 11

Author SHA1 Message Date
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
8 changed files with 1006 additions and 39 deletions
+44 -3
View File
@@ -218,6 +218,44 @@ upstreams:
рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ
не трогаем. Без флажка (`off`) ответ идёт байт-в-байт как раньше. не трогаем. Без флажка (`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
```
### Пример сборки тела (многослойный `patch`) ### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -296,9 +334,12 @@ data class ProviderConf(
val id: String, val id: String,
val url: String, val url: String,
val key: String = "", val key: String = "",
val max_concurrency: Int? = null, // лимит по умолчанию для апстримов провайдера val max_concurrency: Int? = null,
val patch: JsonObject? = null, // ко всем запросам провайдера val patch: JsonObject? = null,
val session_header: String? = null, // заголовок-сессия, считается из истории val session_header: String? = null,
val think_tags: String? = null,
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
) )
@Serializable @Serializable
+45 -2
View File
@@ -13,10 +13,13 @@
| Файл | Что проверяет | | Файл | Что проверяет |
| --- | --- | | --- | --- |
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; | | `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
| `ConfigLogicTest` | Разбор конфига, приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; | | `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; | | `ThinkTagTransformTest` | Non-stream путь `transformThinkMessage`: перенос рассуждений в `reasoning_content`, дописывание к уже имеющемуся, strip, незакрытый блок, отсутствие изменений → null; |
| `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; | | `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; |
| `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. | | `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. |
| `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. |
## Проверка качества тестов (мутационная приёмка) ## Проверка качества тестов (мутационная приёмка)
@@ -36,10 +39,50 @@
- `transformThinkChunk` возвращает null → падают тесты чанков и стрима; - `transformThinkChunk` возвращает null → падают тесты чанков и стрима;
- блок финиш-чанка в `streamSseWithThinkTags` не выполняется → падает тест про удержанный хвост; - блок финиш-чанка в `streamSseWithThinkTags` не выполняется → падает тест про удержанный хвост;
- `holdableSuffix` всегда 0 (хвост тега не удерживается) → падают тесты автомата, чанков и стрима; - `holdableSuffix` всегда 0 (хвост тега не удерживается) → падают тесты автомата, чанков и стрима;
- `strip` начинает отдавать рассуждения → падает тест автомата. - `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. Реальный апстрим проверяется только после деплоя. Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
Правка поля рассуждений — исключение: она проверена живым инстансом против настоящего Console Go (таблица выше) до релиза.
+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** не трогаем: фикс на нашем хопе лечит всех клиентов сразу.
+188 -34
View File
@@ -176,7 +176,8 @@ private suspend fun handleChat(
continue 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) { val forwarded = if (clientWantsStream) {
patched patched
} else { } else {
@@ -234,6 +235,15 @@ private suspend fun handleChat(
failover = true failover = true
return@execute return@execute
} }
if (upstreamStatus in 400..499) {
// 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть.
val errorBody = runCatching { resp.body<String>() }.getOrDefault("")
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" }
val ct = resp.headers["Content-Type"] ?: "application/json"
call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus))
responded = true
return@execute
}
responded = true responded = true
if (clientWantsStream) { if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream" val ct = resp.headers["Content-Type"] ?: "text/event-stream"
@@ -242,16 +252,8 @@ private suspend fun handleChat(
call.respondBytesWriter(ContentType.parse(ct), status) { call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>() val ch = resp.body<ByteReadChannel>()
if (thinkMode == "off") { if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough как раньше. // Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
val buf = ByteArray(8192) streamRawWithDoneContract(ch)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
}
}
} else { } else {
streamSseWithThinkTags(ch, thinkMode) streamSseWithThinkTags(ch, thinkMode)
} }
@@ -292,7 +294,7 @@ private suspend fun handleChat(
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" } log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
throw e throw e
} catch (e: Exception) { } catch (e: Exception) {
log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" } log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" }
failed.add(up.id) failed.add(up.id)
continue continue
} finally { } finally {
@@ -546,27 +548,27 @@ internal fun rebuildFromChunks(sse: String): String {
val data = line.removePrefix("data:").trim() val data = line.removePrefix("data:").trim()
if (data.isEmpty() || data == "[DONE]") return@forEach if (data.isEmpty() || data == "[DONE]") return@forEach
val obj = runCatching { json.parseToJsonElement(data).jsonObject }.getOrNull() ?: return@forEach val obj = runCatching { json.parseToJsonElement(data).jsonObject }.getOrNull() ?: return@forEach
if (id.isEmpty()) id = obj["id"]?.jsonPrimitive?.content ?: "" if (id.isEmpty()) id = (obj["id"] as? JsonPrimitive)?.content ?: ""
if (created == null) created = obj["created"]?.jsonPrimitive?.content?.toLongOrNull() if (created == null) created = (obj["created"] as? JsonPrimitive)?.content?.toLongOrNull()
if (model.isEmpty()) model = obj["model"]?.jsonPrimitive?.content ?: "" if (model.isEmpty()) model = (obj["model"] as? JsonPrimitive)?.content ?: ""
if (systemFingerprint == null) systemFingerprint = obj["system_fingerprint"]?.jsonPrimitive?.content if (systemFingerprint == null) systemFingerprint = (obj["system_fingerprint"] as? JsonPrimitive)?.content
if (serviceTier == null) serviceTier = obj["service_tier"]?.jsonPrimitive?.content if (serviceTier == null) serviceTier = (obj["service_tier"] as? JsonPrimitive)?.content
if (provider == null) provider = obj["provider"] if (provider == null) provider = obj["provider"]
(obj["error"] as? JsonObject)?.let { error = it } (obj["error"] as? JsonObject)?.let { error = it }
(obj["usage"] as? JsonObject)?.let { usage = 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) { for (ch in chArr) {
val c = ch.jsonObject val c = ch as? JsonObject ?: continue
val idx = c["index"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0 val idx = (c["index"] as? JsonPrimitive)?.content?.toIntOrNull() ?: 0
val mc = choices.getOrPut(idx) { MutableChoice() } val mc = choices.getOrPut(idx) { MutableChoice() }
val delta = c["delta"]?.jsonObject val delta = c["delta"] as? JsonObject
if (delta != null) { if (delta != null) {
if (mc.role == null) mc.role = delta["role"]?.jsonPrimitive?.content if (mc.role == null) mc.role = (delta["role"] as? JsonPrimitive)?.content
delta["content"]?.jsonPrimitive?.content?.takeIf { it != "null" }?.let { mc.content.append(it) } (delta["content"] as? JsonPrimitive)?.content?.takeIf { it != "null" }?.let { mc.content.append(it) }
delta["reasoning_content"]?.jsonPrimitive?.content?.takeIf { it != "null" }?.let { mc.reasoning.append(it) } (delta["reasoning_content"] as? JsonPrimitive)?.content?.takeIf { it != "null" }?.let { mc.reasoning.append(it) }
delta["tool_calls"]?.jsonArray?.forEach { tc -> (tc as? JsonObject)?.let { mc.toolCalls.add(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 } c["logprobs"]?.let { mc.logprobs = it }
} }
} }
@@ -618,12 +620,12 @@ internal fun rebuildFromChunks(sse: String): String {
* изменился → null (отдать исходную строку как есть). * изменился → null (отдать исходную строку как есть).
*/ */
internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String? { 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" val addReasoning = thinkMode == "split"
var changed = false var changed = false
val newChoices = choices.map { choiceEl -> val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject val choice = choiceEl as? JsonObject ?: return@map choiceEl
val message = choice["message"]?.jsonObject ?: return@map choiceEl val message = choice["message"] as? JsonObject ?: return@map choiceEl
val contentStr = (message["content"] as? JsonPrimitive)?.takeIf { it.isString }?.content val contentStr = (message["content"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl if (contentStr == null) return@map choiceEl
val splitter = ThinkTagSplitter(thinkMode) val splitter = ThinkTagSplitter(thinkMode)
@@ -644,6 +646,144 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
return JsonObject(obj.toMutableMap().apply { this["choices"] = JsonArray(newChoices) }).toString() 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-тегов. Строки, не начинающиеся * Построчный разбор SSE-стрима с рассечением think-тегов. Строки, не начинающиеся
* с `data:`, и `data: [DONE]` уходят клиенту без изменений (с `\n`). Прочие * с `data:`, и `data: [DONE]` уходят клиенту без изменений (с `\n`). Прочие
@@ -655,15 +795,19 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) { internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>() val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split" val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
while (true) { while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break val line = source.readLine(LineEnding.Lenient) ?: break
when { when {
line.startsWith("data:") -> { line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim() val payload = line.removePrefix("data:").trim()
if (payload == "[DONE]") { if (payload == "[DONE]") {
emitUtf8("data: [DONE]\n\n") sawDone = true
emitUtf8(SSE_DONE_MARKER)
} else { } else {
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull() val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
if (obj != null && hasFinishReason(obj)) sawFinishReason = true
val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) } val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) }
if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n") if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n")
} }
@@ -699,6 +843,12 @@ internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadCha
flush() flush()
} }
} }
// После сброса хвостов дописываем маркер, только если апстрим его не прислал,
// а `finish_reason` в потоке был (при обрыве без него маркер НЕ дописываем).
if (sawFinishReason && !sawDone) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
} }
/** /**
@@ -716,14 +866,14 @@ internal fun transformThinkChunk(
thinkMode: String, thinkMode: String,
addReasoning: Boolean, addReasoning: Boolean,
): String? { ): String? {
val choices = obj["choices"]?.jsonArray ?: return null val choices = obj["choices"] as? JsonArray ?: return null
var changed = false var changed = false
val newChoices = choices.map { choiceEl -> val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject val choice = choiceEl as? JsonObject ?: return@map choiceEl
val delta = choice["delta"]?.jsonObject val delta = choice["delta"] as? JsonObject
val contentStr = (delta?.get("content") as? JsonPrimitive)?.takeIf { it.isString }?.content val contentStr = (delta?.get("content") as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl 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 splitter = splitters.getOrPut(idx) { ThinkTagSplitter(thinkMode) }
val (newContent, reasoning) = splitter.feed(contentStr) val (newContent, reasoning) = splitter.feed(contentStr)
changed = true changed = true
@@ -751,6 +901,8 @@ data class ProviderConf(
val patch: JsonObject? = null, val patch: JsonObject? = null,
val session_header: String? = null, val session_header: String? = null,
val think_tags: String? = null, val think_tags: String? = null,
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
) )
data class UpstreamConf( data class UpstreamConf(
@@ -815,6 +967,8 @@ internal fun parseConfig(root: YamlElement): Config {
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
session_header = m.strOrNull("session_header"), session_header = m.strOrNull("session_header"),
think_tags = m.strOrNull("think_tags"), think_tags = m.strOrNull("think_tags"),
reasoning_field = m.strOrNull("reasoning_field"),
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
) )
} }
@@ -546,4 +546,41 @@ class ConfigLogicTest {
assertEquals("off", effectiveThinkTags(upNull, prov(null))) assertEquals("off", effectiveThinkTags(upNull, prov(null)))
assertEquals("off", effectiveThinkTags(upNull, 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,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), "маркер не должен дублироваться")
}
}