3 Commits
8 ... 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
subochev 334015a2bd test: покрытие think_tags — non-stream, SSE-чанки, обвязка стрима
Build LLM Proxy / Build and push (release) Successful in 44s
- transformThinkMessage/transformThinkChunk/streamSseWithThinkTags сделаны internal
  ради тестируемости (логика не менялась)
- +24 теста: ThinkTagTransformTest (9), ThinkTagChunkTest (8), ThinkTagStreamTest (7)
- ConfigLogicTest: ассерт приоритета источников был неразличающим (провайдер и
  апстрим давали одинаковый результат) — заменён на различающиеся значения
- kotlinx-coroutines-test для тестов каналов ktor
- Gitea Actions: шаг Run tests (jvmTest) — раньше CI тесты не гонял вовсе
- TESTING.md: что покрыто + мутационная приёмка
2026-09-11 17:40:06 +03:00
14 changed files with 1590 additions and 40 deletions
+4
View File
@@ -9,6 +9,10 @@ jobs:
steps:
- name: 'Checkout'
uses: https://github.com/actions/checkout@v4
- name: 'Run tests'
uses: https://git.binom.pw/subochev/devops/build-gradle@main
with:
target: jvmTest
- name: 'Build jar'
uses: https://git.binom.pw/subochev/devops/build-gradle@main
with:
+44 -3
View File
@@ -218,6 +218,44 @@ upstreams:
рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ
не трогаем. Без флажка (`off`) ответ идёт байт-в-байт как раньше.
### Нативное поле рассуждений (`reasoning_field`)
Некоторые шлюзы-апстримы (например, **Console Go / deepseek в thinking-режиме**)
в thinking-режиме **требуют** вернуть им нативное поле рассуждений в каждом
assistant-сообщении с непустым `tool_calls`. Если его нет — апстрим отвечает
`400` (`The `reasoning_content` in the thinking mode must be passed back to the
API.`). Клиенты при этом рассуждения держат в своих форматах: `reasoning`
(строка) и/или `reasoning_details` (массив `{type:"reasoning.text", text, ...}`)
— а нативное поле `reasoning_content` в истории могут и не передавать.
Поля (только на уровне **провайдера**, это свойство шлюза, а не модели):
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) |
| `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе |
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в каждом
assistant-сообщении с непустым `tool_calls` — берёт текст из `reasoning`
(если это непустая строка), иначе склеивает `reasoning_details[*].text`
(только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и
записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее
непустое поле не перезаписывается. Ничего при этом не убирается и не
переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте,
поле просто дополняется. Тронуто только assistant-сообщение с непустым
`tool_calls`; сообщения без `tool_calls` и `user`/`tool`/`system` не меняются.
Если текста рассуждений нет вовсе — поле не добавляется, кроме случая
`reasoning_empty_ok: true` (тогда пишется пустая строка `""`). Если
`reasoning_field` не задан — тело не меняется вовсе.
```yaml
providers:
- id: opencode
url: "https://opencode.ai/zen/go/v1"
reasoning_field: reasoning_content # Console Go / deepseek в thinking-режиме
# reasoning_empty_ok: true # опционально; дефолт false
```
### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -296,9 +334,12 @@ data class ProviderConf(
val id: String,
val url: String,
val key: String = "",
val max_concurrency: Int? = null, // лимит по умолчанию для апстримов провайдера
val patch: JsonObject? = null, // ко всем запросам провайдера
val session_header: String? = null, // заголовок-сессия, считается из истории
val max_concurrency: Int? = null,
val patch: JsonObject? = null,
val session_header: String? = null,
val think_tags: String? = null,
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
)
@Serializable
+88
View File
@@ -0,0 +1,88 @@
# Тестирование llm-proxy
Все тесты — обычные модульные, живут в `src/commonTest/kotlin/pw/binom/llmproxy/`.
## Запуск
- `./gradlew jvmTest` — прогнать все тесты;
- `./gradlew clean jvmTest fatJar` — полная сборка с нуля;
- результат смотреть в `build/test-results/jvmTest/*.xml` (атрибуты `tests`/`failures`/`errors`), потому что строки вида «N tests completed» печатаются только при падениях.
## Что покрыто
| Файл | Что проверяет |
| --- | --- |
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
| `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, неприкосновенность сообщений без `tool_calls` (в т.ч. `[]`) и `user`/`tool`, пустая строка при `emptyOk`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
| `NullToleranceTest` | Устойчивость разбора ответов апстрима к JSON-`null` (`choices`, `delta.tool_calls`, `delta.content`) — пропуск вместо исключения; |
| `ThinkTagTransformTest` | Non-stream путь `transformThinkMessage`: перенос рассуждений в `reasoning_content`, дописывание к уже имеющемуся, strip, незакрытый блок, отсутствие изменений → null; |
| `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; |
| `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. |
| `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. |
## Проверка качества тестов (мутационная приёмка)
Приём по шагам:
1. Забэкапить файл.
2. Внести РОВНО одну поломку в боевой код.
3. Прогнать `./gradlew cleanJvmTest jvmTest`.
4. Посмотреть XML — тест, который не упал, считается пустым.
5. Откатить (`git checkout -- <файл>`).
Обязательно: `cleanJvmTest` обязателен, иначе прогон не перезапустится.
Проверенные мутации, каждая из которых ДОЛЖНА ронять тесты:
- `transformThinkMessage` возвращает null → падают тесты non-stream;
- `transformThinkChunk` возвращает null → падают тесты чанков и стрима;
- блок финиш-чанка в `streamSseWithThinkTags` не выполняется → падает тест про удержанный хвост;
- `holdableSuffix` всегда 0 (хвост тега не удерживается) → падают тесты автомата, чанков и стрима;
- `strip` начинает отдавать рассуждения → падает тест автомата;
- `missingDoneMarker` всегда `false` → падают тесты сырого пути (`rawStreamIsByteExactAndAppendsDone`) и тест разрезанного маркера;
- `if (sawFinishReason && !sawDone)` → `if (!sawDone)` в `streamSseWithThinkTags` → падает `thinkStreamTruncatedNeedsNoMarker`;
- `if (sawFinishReason && !sawDone)` → `if (false)` в `streamSseWithThinkTags` → падает `thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker`.
Правило: боевой код нельзя подгонять под тест; если тест не проходит, неверен тест.
## Поле рассуждений для Console Go (релиз 11) — замеры приёмки
Причина правки: апстрим `deepseek-v4.1-flash` у провайдера `opencode` (Console Go) в thinking-режиме
требует `reasoning_content` в assistant-сообщениях с `tool_calls`, а клиент opencode присылает
рассуждения как `reasoning` + `reasoning_details` — отсюда `400 The reasoning_content in the
thinking mode must be passed back to the API`.
**Границы требования (замер прямыми запросами к Console Go, одинаковое тело):**
| assistant-сообщение | HTTP |
| --- | --- |
| с `tool_calls`, без reasoning вовсе | 400 |
| с `tool_calls`, `reasoning` + `reasoning_details` (форма opencode) | 400 |
| с `tool_calls`, только `reasoning_details` | 400 |
| с `tool_calls`, `reasoning_content` + `reasoning` + `reasoning_details` (аддитивно) | 200 |
| с `tool_calls`, `reasoning_content: ""` | 200 |
| без `tool_calls`, без reasoning | 200 |
**Живая приёмка (локальный инстанс на 8101, конфиг-копия боевого, модель с единственным апстримом
Console Go; тело — как у opencode: assistant + `reasoning` + `reasoning_details` + `tool_calls`):**
| Конфиг | `stream=false` | `stream=true` |
| --- | --- | --- |
| без правки (`reasoning_field` не задан) | **400** — та самая ошибка про `reasoning_content` | **400** |
| с правкой (`reasoning_field: reasoning_content`) | **200**, модель продолжила диалог после tool-результата | **200** |
**Регрессия на боевых цепочках (тот же инстанс, то же тело с `tool_calls`):** `codding-big` → 200
(ушло на `minimax-m3`, поле не добавляется — флаг объявлен только у провайдера `opencode`),
`codding` → 200 (`qwen-3.8`, local), `assistant` → 200 (Console Go с правкой).
**Мутационная приёмка:** 4 мутации из ТЗ отработал кодер, две проверены вручную
(`cleanAllTests jvmTest`): снятие охраны `tool_calls` в `applyReasoningField` → падают
`assistantWithoutToolCallsIsUntouched` и `assistantWithEmptyToolCallsIsUntouched`;
возврат `obj["choices"]?.jsonArray` в `rebuildFromChunks` → падает `rebuildFromChunksToleratesNullChoices`.
## Известное ограничение
Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
Правка поля рассуждений — исключение: она проверена живым инстансом против настоящего Console Go (таблица выше) до релиза.
+1
View File
@@ -31,6 +31,7 @@ kotlin {
val commonTest by getting {
dependencies {
implementation(libs.kotlin.test)
implementation(libs.kotlinx.coroutines.test)
}
}
val jvmMain by getting {
+280
View File
@@ -0,0 +1,280 @@
# ТЗ: выравнивание контракта конца SSE (дописать `data: [DONE]`)
Задача-источник: Vikunja **#71** («в opencode модель додумала, идёт ожидание ~45 с»).
## 1. Что происходит сейчас (проверено живыми замерами 11.09.2026)
Цепочка: opencode → `llm.binom.pw` (Bifrost v1.5.0) → `llm-proxy:8100` → `api.minimax.io`.
- **MiniMax не присылает терминатор `data: [DONE]`.** После чанка с `finish_reason: "stop"` он просто закрывает
сокет (сразу). Проверено и с `reasoning_split`, и с `stream_options.include_usage`.
- **llm-proxy отдаёт поток корректно**, но клиент ждёт маркер: при `think_tags: off` (это как раз MiniMax-апстрим
`minimax-m3` модели `codding-big`) работает сырой байтовый passthrough — что прислал апстрим, то и ушло.
Маркера нет → нет маркера.
- **Bifrost v1.5.0** для custom-провайдеров считает `[DONE]` обязательным (`ProviderSendsDoneMarker` = true;
флага `custom_provider_config.does_not_send_done_marker` в этой версии нет). Маркера нет → его цикл чтения
не завершается, и стрим закрывается только когда умирает keep-alive-сокет llm-proxy — через **45.00 с**
(Ktor CIO `connectionIdleTimeoutSeconds = 45`). После этого Bifrost дописывает клиенту `[DONE]` и закрывает поток.
Отсюда «модель додумала, а opencode ещё думает».
- Контроль: тот же Bifrost и та же тула, апстрим со своим `[DONE]` (RouterAI `local/deepseek/deepseek-v4-flash-0731`)
— закрытие за **0.00 с**.
## 2. Требуемое поведение (контракт)
Стрим, который отдаёт llm-proxy клиенту, **всегда** заканчивается маркером `data: [DONE]\n\n`, если апстрим
завершил поток штатно (в потоке был непустой `finish_reason`), — даже когда сам апстрим маркер не прислал.
Правило ровно в двух частях:
1. Апстрим закрыл поток, и в потоке был непустой `finish_reason`, но `data: [DONE]` мы не видели →
дописать клиенту `data: [DONE]\n\n` и `flush()`.
2. Апстрим оборвался, **не** прислав `finish_reason` (усечённый/битый поток) → маркер **НЕ** дописывать,
отдать как есть.
Пункт 2 — не забыть и не «улучшить»: именно по отсутствию маркера Bifrost отличает усечённый стрим от
успешного (его детект `SSEStreamEndedOnMarker`). Если дописать маркер при обрыве без `finish_reason`,
мы превратим обрыв в «успех» и сломаем диагностику.
Обязательно: **байты уже существующих чанков не меняются** (passthrough остаётся byte-exact), маркер
дописывается только В КОНЦЕ, только один раз, и никогда не дублируется, если апстрим прислал свой.
## 3. Правки в коде
Файл **один**: `src/commonMain/kotlin/pw/binom/llmproxy/Main.kt`.
### 3.1. Константы и хелперы (добавить рядом с блоком think-хелперов)
```kotlin
/** Терминальный маркер SSE, который клиенты (в т.ч. Bifrost) ждут как конец потока. */
private const val SSE_DONE_MARKER = "data: [DONE]\n\n"
/** Размер скользящего окна детектора: маркер может разрезаться границей чтения. */
private const val STREAM_SCAN_WINDOW = 128
/**
* Есть ли в тексте `finish_reason` со значением, отличным от `null`. Экранированные
* вхождения (`\"finish_reason\\\":\\\"stop\\\"` внутри содержимого ответа) не считаются:
* смотрим только на неэкранированную кавычку.
*/
internal fun hasNonNullFinishReasonText(text: String): Boolean {
var from = 0
while (true) {
val i = text.indexOf("finish_reason", from)
if (i < 0) return false
val quote = i - 1
val escaped = quote >= 0 && text[quote] == '"' && quote > 0 && text[quote - 1] == '\\'
if (!escaped) {
val rest = text.substring(i + "finish_reason".length)
.dropWhile { it == ' ' || it == ':' || it == '"' }
if (!rest.startsWith("null")) return true
}
from = i + 1
}
}
/** Структурная проверка чанка: в `choices[*].finish_reason` есть непустая строка. */
internal fun hasFinishReason(obj: JsonObject): Boolean =
obj["choices"]?.jsonArray?.any { el ->
val fr = (el.jsonObject["finish_reason"] as? JsonPrimitive)?.takeIf { it.isString }?.content
!fr.isNullOrEmpty()
} ?: false
/**
* Детектор контракта конца SSE для сырого passthrough-потока (`think_tags: off`).
* Копит скользящее окно последних байт, чтобы маркер `finish_reason`/`[DONE]`,
* разрезанный границей чтения, всё равно был распознан, и запоминает два факта:
* видели ли непустой `finish_reason` и видели ли `[DONE]`.
*/
internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDOW) {
var sawFinishReason = false
private set
var sawDone = false
private set
private var window: String = ""
fun feed(bytes: ByteArray, from: Int = 0, length: Int = bytes.size) {
if (length <= 0) return
val scan = window + bytes.decodeToString(from, from + length)
if (!sawDone && scan.contains("[DONE]")) sawDone = true
if (!sawFinishReason && hasNonNullFinishReasonText(scan)) sawFinishReason = true
window = if (scan.length > windowSize) scan.substring(scan.length - windowSize) else scan
}
/**
* Маркер нужен, только если апстрим закрыл поток после штатного `finish_reason`,
* но сам `[DONE]` не прислал. При обрыве БЕЗ `finish_reason` маркер НЕ дописываем.
*/
val missingDoneMarker: Boolean get() = sawFinishReason && !sawDone
}
/**
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
*/
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
val detector = StreamEndDetector()
val buf = ByteArray(8192)
while (true) {
val n = source.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
detector.feed(buf, 0, n)
}
}
if (detector.missingDoneMarker) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
```
Импорты под новые типы (`ByteReadChannel`, `ByteWriteChannel`, `readAvailable`, `writeFully`, `flush`,
`jsonArray`, `jsonObject`, `JsonPrimitive`) — проверить, что они уже есть в файле; добавлять только недостающие.
### 3.2. Правка ветки `think_tags: off` в `handleChat`
Было (внутри `call.respondBytesWriter(...)`):
```kotlin
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough как раньше.
val buf = ByteArray(8192)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
}
}
} else {
streamSseWithThinkTags(ch, thinkMode)
}
```
Стало:
```kotlin
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
ch.streamRawWithDoneContract(this)
} else {
streamSseWithThinkTags(ch, thinkMode)
}
```
(если `this` в этом контексте — не `ByteWriteChannel`, использовать тот ресивер, который уже использовался
существующими вызовами `writeFully`/`flush` в этой лямбде; компилятор — источник истины).
### 3.3. Правка `streamSseWithThinkTags`
Добавить два флага и запись фактов:
```kotlin
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break
when {
line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim()
if (payload == "[DONE]") {
sawDone = true
emitUtf8(SSE_DONE_MARKER)
} else {
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
if (obj != null && hasFinishReason(obj)) sawFinishReason = true
val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) }
if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n")
}
}
else -> emitUtf8("$line\n")
}
flush()
}
// ... существующий блок сброса удержанных хвостов сплиттеров — БЕЗ изменений ...
// НОВОЕ: после хвостов дописать маркер, если апстрим его не прислал, а finish_reason был.
if (sawFinishReason && !sawDone) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
```
Порядок важен: маркер дописывается **после** блока сброса хвостов (сначала текст, потом терминатор).
Остальную логику функции (`transformThinkChunk`, `emitUtf8("$line\n")` для служебных строк) не менять.
## 4. Тесты (новый файл)
`src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt` — по образцу существующего
`ThinkTagStreamTest` (тот же приём: `ByteChannel(autoFlush = true)`, запись входа, `close(null)`,
чтение выхода в буфер; для сырого потока вызывать `streamRawWithDoneContract`).
Обязательные случаи (имена — как ниже, чтобы можно было ссылаться в отчёте):
| Тест | Вход | Ожидание |
| --- | --- | --- |
| `textHelperSeesFinishReason` | строка `{"choices":[{"finish_reason":"stop"}]}` | `hasNonNullFinishReasonText` = true |
| `textHelperIgnoresNullFinishReason` | `{"choices":[{"finish_reason":null}]}` | false |
| `textHelperIgnoresEscapedOccurrenceInContent` | `{"delta":{"content":"\"finish_reason\":\"stop\""}}` | false |
| `detectorSeesFinishReasonSplitAcrossReads` | два `feed`: `...finish_rea` + `son":"stop"...` | `sawFinishReason` = true, `missingDoneMarker` = true |
| `detectorSeesDoneSplitAcrossReads` | два `feed`: `data: [DO` + `NE]` | `sawDone` = true, `missingDoneMarker` = false |
| `detectorSeesDoneAndFinishReason` | `finish_reason` + `[DONE]` в одном feed | `missingDoneMarker` = false |
| `detectorTruncatedStreamNeedsNoMarker` | `feed` с контентом без `finish_reason` | `missingDoneMarker` = false |
| `rawStreamIsByteExactAndAppendsDone` | SSE с `finish_reason`, без `[DONE]` | выход = вход + `"data: [DONE]\n\n"` (байт-в-байт префикс), маркер ровно один |
| `rawStreamTruncatedPassesThroughUntouched` | SSE без `finish_reason` и без `[DONE]` | выход = вход (ни одного `[DONE]`) |
| `rawStreamKeepsUpstreamDoneMarker` | SSE с `finish_reason` и своим `[DONE]` | `[DONE]` ровно один (не дублируется) |
| `thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker` | `streamSseWithThinkTags(..., "split")`, чанк с `finish_reason:"stop"`, без `[DONE]` | выход заканчивается `data: [DONE]\n\n`, ровно один маркер |
| `thinkStreamTruncatedNeedsNoMarker` | mode `split`, чанк с контентом, без `finish_reason` | в выходе нет `[DONE]` |
| `thinkStreamKeepsSingleDone` | mode `split`, `finish_reason` + `data: [DONE]` | ровно один `[DONE]` |
## 5. Мутационная приёмка (по TESTING.md)
После того как тесты зелёные, проверить их силу — по одному изменению, каждый раз с
`./gradlew cleanJvmTest jvmTest`, затем `git checkout --` на файл:
1. `missingDoneMarker` → всегда `false` (детектор сырого пути) — должны упасть тесты сырого пути
(`rawStreamIsByteExactAndAppendsDone`) и тест разрезанного маркера;
2. `if (sawFinishReason && !sawDone)` → `if (!sawDone)` в `streamSseWithThinkTags` — должен упасть
`thinkStreamTruncatedNeedsNoMarker` (это защита от «обрыв выглядит как успех»);
3. `if (sawFinishReason && !sawDone)` → `if (false)` — должен упасть
`thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker`.
Тест, который при такой поломке остаётся зелёным, считается пустым.
## 6. Команды проверки
```bash
./gradlew jvmTest # все тесты зелёные; сейчас в проекте 66 тестов, стало больше
./gradlew cleanJvmTest jvmTest # прогон с нуля (обязателен для мутационной приёмки)
./gradlew clean jvmTest fatJar # полная сборка
```
Числа тестов смотреть в `build/test-results/jvmTest/*.xml` (атрибуты `tests` / `failures` / `errors`),
а не по строкам в консоли.
## 7. Границы работ
- Трогать **только**: `src/commonMain/kotlin/pw/binom/llmproxy/Main.kt`,
`src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt`, `TESTING.md` (добавить строку про
новый тест-файл и пункты мутационной приёмки).
- **НЕ** менять `ThinkTagSplitter.kt`, non-stream путь, `transformThinkMessage`, `transformThinkChunk`,
фрейминг существующих чанков, `README.md`/`CONFIG.md`.
- **НЕ** добавлять зависимости, не менять версии в `build.gradle.kts`, не менять `config.yaml`.
- **НЕ** коммитить и не пушить.
- Внешние сети/апстримы не дёргать — только локальные тесты.
## 8. Вне области (и почему)
- **Обрыв без `finish_reason` остаётся без маркера** — так усечённый стрим не выглядит успешным
(см. п. 2). Да, в этом случае Bifrost снова будет ждать закрытия сокета; это осознанный выбор
в пользу честности контракта, а не «ускорения любой ценой».
- **Non-stream путь** не трогаем: там тула сама собирает ответ из чанков (`rebuildFromChunks`) и
маркер не нужен.
- **Bifrost и сам MiniMax** не трогаем: фикс на нашем хопе лечит всех клиентов сразу.
+1
View File
@@ -19,6 +19,7 @@ kotlinx-serialization-json = { module = "org.jetbrains.kotlinx:kotlinx-serializa
yamlkt = { module = "net.mamoe.yamlkt:yamlkt", version.ref = "yamlkt" }
kotlinx-coroutines-core = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-core", version.ref = "coroutines" }
kotlinx-datetime = { module = "org.jetbrains.kotlinx:kotlinx-datetime", version.ref = "datetime" }
kotlinx-coroutines-test = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-test", version.ref = "coroutines" }
kotlinx-io-core = { module = "org.jetbrains.kotlinx:kotlinx-io-core", version.ref = "kotlinxIo" }
logback-classic = { module = "ch.qos.logback:logback-classic", version.ref = "logback" }
kotlin-test = { module = "org.jetbrains.kotlin:kotlin-test" }
+190 -36
View File
@@ -176,7 +176,8 @@ private suspend fun handleChat(
continue
}
val patched = buildBody(bodyJson, provider, up, modelConf)
val patched0 = buildBody(bodyJson, provider, up, modelConf)
val patched = applyReasoningField(patched0, provider.reasoning_field, provider.reasoning_empty_ok)
val forwarded = if (clientWantsStream) {
patched
} else {
@@ -234,6 +235,15 @@ private suspend fun handleChat(
failover = true
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
if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream"
@@ -242,16 +252,8 @@ private suspend fun handleChat(
call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>()
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough как раньше.
val buf = ByteArray(8192)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
}
}
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
streamRawWithDoneContract(ch)
} else {
streamSseWithThinkTags(ch, thinkMode)
}
@@ -292,7 +294,7 @@ private suspend fun handleChat(
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
throw e
} catch (e: Exception) {
log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" }
log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" }
failed.add(up.id)
continue
} finally {
@@ -546,27 +548,27 @@ internal fun rebuildFromChunks(sse: String): String {
val data = line.removePrefix("data:").trim()
if (data.isEmpty() || data == "[DONE]") return@forEach
val obj = runCatching { json.parseToJsonElement(data).jsonObject }.getOrNull() ?: return@forEach
if (id.isEmpty()) id = obj["id"]?.jsonPrimitive?.content ?: ""
if (created == null) created = obj["created"]?.jsonPrimitive?.content?.toLongOrNull()
if (model.isEmpty()) model = obj["model"]?.jsonPrimitive?.content ?: ""
if (systemFingerprint == null) systemFingerprint = obj["system_fingerprint"]?.jsonPrimitive?.content
if (serviceTier == null) serviceTier = obj["service_tier"]?.jsonPrimitive?.content
if (id.isEmpty()) id = (obj["id"] as? JsonPrimitive)?.content ?: ""
if (created == null) created = (obj["created"] as? JsonPrimitive)?.content?.toLongOrNull()
if (model.isEmpty()) model = (obj["model"] as? JsonPrimitive)?.content ?: ""
if (systemFingerprint == null) systemFingerprint = (obj["system_fingerprint"] as? JsonPrimitive)?.content
if (serviceTier == null) serviceTier = (obj["service_tier"] as? JsonPrimitive)?.content
if (provider == null) provider = obj["provider"]
(obj["error"] as? JsonObject)?.let { error = it }
(obj["usage"] as? JsonObject)?.let { usage = it }
val chArr = obj["choices"]?.jsonArray ?: return@forEach
val chArr = obj["choices"] as? JsonArray ?: return@forEach
for (ch in chArr) {
val c = ch.jsonObject
val idx = c["index"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0
val c = ch as? JsonObject ?: continue
val idx = (c["index"] as? JsonPrimitive)?.content?.toIntOrNull() ?: 0
val mc = choices.getOrPut(idx) { MutableChoice() }
val delta = c["delta"]?.jsonObject
val delta = c["delta"] as? JsonObject
if (delta != null) {
if (mc.role == null) mc.role = delta["role"]?.jsonPrimitive?.content
delta["content"]?.jsonPrimitive?.content?.takeIf { it != "null" }?.let { mc.content.append(it) }
delta["reasoning_content"]?.jsonPrimitive?.content?.takeIf { it != "null" }?.let { mc.reasoning.append(it) }
delta["tool_calls"]?.jsonArray?.forEach { tc -> (tc as? JsonObject)?.let { mc.toolCalls.add(it) } }
if (mc.role == null) mc.role = (delta["role"] as? JsonPrimitive)?.content
(delta["content"] as? JsonPrimitive)?.content?.takeIf { it != "null" }?.let { mc.content.append(it) }
(delta["reasoning_content"] as? JsonPrimitive)?.content?.takeIf { it != "null" }?.let { mc.reasoning.append(it) }
(delta["tool_calls"] as? JsonArray)?.forEach { tc -> (tc as? JsonObject)?.let { mc.toolCalls.add(it) } }
}
c["finish_reason"]?.jsonPrimitive?.content?.takeIf { it.isNotEmpty() && it != "null" }?.let { mc.finishReason = it }
(c["finish_reason"] as? JsonPrimitive)?.content?.takeIf { it.isNotEmpty() && it != "null" }?.let { mc.finishReason = it }
c["logprobs"]?.let { mc.logprobs = it }
}
}
@@ -618,12 +620,12 @@ internal fun rebuildFromChunks(sse: String): String {
* изменился → null (отдать исходную строку как есть).
*/
internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String? {
val choices = obj["choices"]?.jsonArray ?: return null
val choices = obj["choices"] as? JsonArray ?: return null
val addReasoning = thinkMode == "split"
var changed = false
val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject
val message = choice["message"]?.jsonObject ?: return@map choiceEl
val choice = choiceEl as? JsonObject ?: return@map choiceEl
val message = choice["message"] as? JsonObject ?: return@map choiceEl
val contentStr = (message["content"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl
val splitter = ThinkTagSplitter(thinkMode)
@@ -644,6 +646,144 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
return JsonObject(obj.toMutableMap().apply { this["choices"] = JsonArray(newChoices) }).toString()
}
/**
* Достроить нативное поле рассуждений для апстримов, которые его требуют
* (Console Go / deepseek в thinking-режиме): если у провайдера объявлено
* `reasoningField`, то в каждом assistant-сообщении с непустым `tool_calls`
* добавляем это поле, ЕСЛИ его там ещё нет. Текст берём из `reasoning`
* (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`).
* Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста
* пишем пустую строку, только если `emptyOk`.
* Тело возвращается без изменений (тот же объект), если менять нечего.
*/
internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Boolean): JsonObject {
if (field == null) return body
val messages = body["messages"] as? JsonArray ?: return body
var changed = false
val newMessages = messages.map { el ->
val msg = el as? JsonObject ?: return@map el
val role = (msg["role"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (role != "assistant") return@map el
val toolCalls = msg["tool_calls"] as? JsonArray
if (toolCalls == null || toolCalls.isEmpty()) return@map el
val existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (existing != null && existing.isNotEmpty()) return@map el
val text = reasoningTextOf(msg)
if (text.isEmpty() && !emptyOk) return@map el
changed = true
JsonObject(msg.toMutableMap().apply { this[field] = JsonPrimitive(text) })
}
if (!changed) return body
return JsonObject(body.toMutableMap().apply { this["messages"] = JsonArray(newMessages) })
}
/**
* Текст рассуждений сообщения: `reasoning` (если непустая строка), иначе
* склейка `reasoning_details[*].text` через "\n" — только элементы, у которых
* `type` отсутствует или равен "reasoning.text".
*/
internal fun reasoningTextOf(msg: JsonObject): String {
val reasoning = (msg["reasoning"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (reasoning != null && reasoning.isNotEmpty()) return reasoning
val details = msg["reasoning_details"] as? JsonArray ?: return ""
val parts = details.mapNotNull { el ->
val d = el as? JsonObject ?: return@mapNotNull null
val type = (d["type"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (type != null && type != "reasoning.text") return@mapNotNull null
(d["text"] as? JsonPrimitive)?.takeIf { it.isString }?.content
}
return parts.joinToString("\n")
}
/** Терминальный маркер SSE, который клиенты (в т.ч. Bifrost) ждут как конец потока. */
private const val SSE_DONE_MARKER = "data: [DONE]\n\n"
/** Размер скользящего окна детектора: маркер может разрезаться границей чтения. */
private const val STREAM_SCAN_WINDOW = 128
/**
* Есть ли в тексте `finish_reason` со значением, отличным от `null`. Экранированные
* вхождения (`\"finish_reason\\\":\\\"stop\\\"` внутри содержимого ответа) не считаются:
* смотрим только на неэкранированную кавычку.
*/
internal fun hasNonNullFinishReasonText(text: String): Boolean {
var from = 0
while (true) {
val i = text.indexOf("finish_reason", from)
if (i < 0) return false
val quote = i - 1
val escaped = quote >= 0 && text[quote] == '"' && quote > 0 && text[quote - 1] == '\\'
if (!escaped) {
val rest = text.substring(i + "finish_reason".length)
.dropWhile { it == ' ' || it == ':' || it == '"' }
if (!rest.startsWith("null")) return true
}
from = i + 1
}
}
/** Структурная проверка чанка: в `choices[*].finish_reason` есть непустая строка. */
internal fun hasFinishReason(obj: JsonObject): Boolean =
(obj["choices"] as? JsonArray)?.any { el ->
(el as? JsonObject)?.let { choice ->
val fr = (choice["finish_reason"] as? JsonPrimitive)?.takeIf { it.isString }?.content
!fr.isNullOrEmpty()
} ?: false
} ?: false
/**
* Детектор контракта конца SSE для сырого passthrough-потока (`think_tags: off`).
* Копит скользящее окно последних байт, чтобы маркер `finish_reason`/`[DONE]`,
* разрезанный границей чтения, всё равно был распознан, и запоминает два факта:
* видели ли непустой `finish_reason` и видели ли `[DONE]`.
*/
internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDOW) {
var sawFinishReason = false
private set
var sawDone = false
private set
private var window: String = ""
fun feed(bytes: ByteArray, from: Int = 0, length: Int = bytes.size) {
if (length <= 0) return
val scan = window + bytes.decodeToString(from, from + length)
if (!sawDone && scan.contains("[DONE]")) sawDone = true
if (!sawFinishReason && hasNonNullFinishReasonText(scan)) sawFinishReason = true
window = if (scan.length > windowSize) scan.substring(scan.length - windowSize) else scan
}
/**
* Маркер нужен, только если апстрим закрыл поток после штатного `finish_reason`,
* но сам `[DONE]` не прислал. При обрыве БЕЗ `finish_reason` маркер НЕ дописываем.
*/
val missingDoneMarker: Boolean get() = sawFinishReason && !sawDone
}
/**
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
*/
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
val detector = StreamEndDetector()
val buf = ByteArray(8192)
while (true) {
val n = source.readAvailable(buf)
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
flush()
detector.feed(buf, 0, n)
}
}
if (detector.missingDoneMarker) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
/**
* Построчный разбор SSE-стрима с рассечением think-тегов. Строки, не начинающиеся
* с `data:`, и `data: [DONE]` уходят клиенту без изменений (с `\n`). Прочие
@@ -652,18 +792,22 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
*/
private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break
when {
line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim()
if (payload == "[DONE]") {
emitUtf8("data: [DONE]\n\n")
sawDone = true
emitUtf8(SSE_DONE_MARKER)
} else {
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
if (obj != null && hasFinishReason(obj)) sawFinishReason = true
val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) }
if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n")
}
@@ -699,6 +843,12 @@ private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChan
flush()
}
}
// После сброса хвостов дописываем маркер, только если апстрим его не прислал,
// а `finish_reason` в потоке был (при обрыве без него маркер НЕ дописываем).
if (sawFinishReason && !sawDone) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
/**
@@ -710,20 +860,20 @@ private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChan
* Чанк без `choices` или без строкового `delta.content` не меняется —
* возвращается null (отдать исходную строку как есть).
*/
private fun transformThinkChunk(
internal fun transformThinkChunk(
obj: JsonObject,
splitters: MutableMap<Int, ThinkTagSplitter>,
thinkMode: String,
addReasoning: Boolean,
): String? {
val choices = obj["choices"]?.jsonArray ?: return null
val choices = obj["choices"] as? JsonArray ?: return null
var changed = false
val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject
val delta = choice["delta"]?.jsonObject
val choice = choiceEl as? JsonObject ?: return@map choiceEl
val delta = choice["delta"] as? JsonObject
val contentStr = (delta?.get("content") as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl
val idx = choice["index"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0
val idx = (choice["index"] as? JsonPrimitive)?.content?.toIntOrNull() ?: 0
val splitter = splitters.getOrPut(idx) { ThinkTagSplitter(thinkMode) }
val (newContent, reasoning) = splitter.feed(contentStr)
changed = true
@@ -751,6 +901,8 @@ data class ProviderConf(
val patch: JsonObject? = null,
val session_header: String? = null,
val think_tags: String? = null,
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
)
data class UpstreamConf(
@@ -815,6 +967,8 @@ internal fun parseConfig(root: YamlElement): Config {
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
session_header = m.strOrNull("session_header"),
think_tags = m.strOrNull("think_tags"),
reasoning_field = m.strOrNull("reasoning_field"),
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
)
}
@@ -533,7 +533,7 @@ class ConfigLogicTest {
// значение у апстрима — берётся оно, провайдер игнорируется
assertEquals("split", effectiveThinkTags(upSplit, null))
assertEquals("split", effectiveThinkTags(upSplit, prov("true")))
assertEquals("split", effectiveThinkTags(upSplit, prov("strip")))
// апстрим null — берётся провайдерский
assertEquals("split", effectiveThinkTags(upNull, prov("split")))
@@ -546,4 +546,41 @@ class ConfigLogicTest {
assertEquals("off", effectiveThinkTags(upNull, prov(null)))
assertEquals("off", effectiveThinkTags(upNull, null))
}
@Test
fun parseConfigReadsReasoningFieldAndEmptyOkOnProvider() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
reasoning_field: reasoning_content
reasoning_empty_ok: true
- id: p2
url: "https://y.ru/api/v1"
reasoning_empty_ok: false
models:
- name: m1
upstreams: []
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals("reasoning_content", cfg.providers[0].reasoning_field)
assertEquals(true, cfg.providers[0].reasoning_empty_ok)
assertEquals(null, cfg.providers[1].reasoning_field)
assertEquals(false, cfg.providers[1].reasoning_empty_ok)
}
@Test
fun parseConfigDefaultsReasoningFieldsWhenAbsent() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
models:
- name: m1
upstreams: []
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals(null, cfg.providers[0].reasoning_field)
assertEquals(false, cfg.providers[0].reasoning_empty_ok)
}
}
@@ -0,0 +1,71 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNull
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class NullToleranceTest {
private val json = Json { ignoreUnknownKeys = true }
@Test
fun rebuildFromChunksToleratesNullChoices() {
// Чанк с choices:null не должен ронять сборку; значимые поля сохраняются.
val sse = """
data: {"id":"c1","created":123,"model":"m","choices":null}
data: {"id":"c1","created":123,"model":"m","choices":[{"index":0,"delta":{"content":"Hi"}}]}
data: [DONE]
""".trimIndent()
val out = json.parseToJsonElement(rebuildFromChunks(sse)).jsonObject
assertEquals("c1", out["id"]?.jsonPrimitive?.content)
assertEquals("m", out["model"]?.jsonPrimitive?.content)
assertEquals(123, out["created"]?.jsonPrimitive?.content?.toLong())
val choice = out["choices"]?.jsonArray?.get(0)?.jsonObject
assertEquals("Hi", choice?.get("message")?.jsonObject?.get("content")?.jsonPrimitive?.content)
}
@Test
fun rebuildFromChunksToleratesNullDeltaToolCalls() {
val sse = """
data: {"id":"c1","model":"m","choices":[{"index":0,"delta":{"role":"assistant","tool_calls":null},"finish_reason":"stop"}]}
data: [DONE]
""".trimIndent()
val out = json.parseToJsonElement(rebuildFromChunks(sse)).jsonObject
val choice = out["choices"]?.jsonArray?.get(0)?.jsonObject
assertEquals("assistant", choice?.get("message")?.jsonObject?.get("role")?.jsonPrimitive?.content)
}
@Test
fun rebuildFromChunksToleratesNullContent() {
val sse = """
data: {"id":"c1","model":"m","choices":[{"index":0,"delta":{"role":"assistant","content":null},"finish_reason":"stop"}]}
data: [DONE]
""".trimIndent()
val out = json.parseToJsonElement(rebuildFromChunks(sse)).jsonObject
val choice = out["choices"]?.jsonArray?.get(0)?.jsonObject
assertEquals("assistant", choice?.get("message")?.jsonObject?.get("role")?.jsonPrimitive?.content)
}
@Test
fun hasFinishReasonToleratesNullChoices() {
val obj = json.parseToJsonElement("""{"choices":null}""").jsonObject
assertFalse(hasFinishReason(obj))
}
@Test
fun transformThinkChunkToleratesNullChoices() {
val obj = json.parseToJsonElement("""{"choices":null}""").jsonObject
assertNull(transformThinkChunk(obj, mutableMapOf(), "split", true))
}
@Test
fun transformThinkMessageToleratesNullChoices() {
val obj = json.parseToJsonElement("""{"choices":null}""").jsonObject
assertNull(transformThinkMessage(obj, "split"))
}
}
@@ -0,0 +1,165 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive
import kotlinx.serialization.json.jsonObject
class ReasoningFieldTest {
private val json = Json { ignoreUnknownKeys = true }
// Минимальный непустой tool_call для assistant-сообщения.
private val TC = """{"id":"t1","type":"function","function":{"name":"f","arguments":"{}"}}"""
private fun str(obj: JsonObject, key: String): String? = (obj[key] as? JsonPrimitive)?.content
private fun outMsgs(out: JsonObject): List<JsonObject> =
(out["messages"] as? JsonArray)?.mapNotNull { it as? JsonObject } ?: emptyList()
// Собрать тело `{"messages":[...]}` и прогнать через applyReasoningField.
private fun apply(messagesJson: String, field: String?, emptyOk: Boolean): JsonObject {
val body = json.parseToJsonElement("""{"messages":[$messagesJson]}""").jsonObject
return applyReasoningField(body, field, emptyOk)
}
@Test
fun documentedOpencodeShapeGetsNativeReasoningContent() {
// Форма клиента opencode: reasoning (строка) + reasoning_details (массив),
// нативного поля нет — добавляем его, текст берём из рассуждений.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning":"R","reasoning_details":[{"type":"reasoning.text","text":"R","index":0}]}"""
val out = apply(m, "reasoning_content", false)
val m0 = outMsgs(out)[0]
assertEquals("R", str(m0, "reasoning_content"))
assertEquals("R", str(m0, "reasoning"))
assertTrue((m0["tool_calls"] as? JsonArray)?.isNotEmpty() == true)
assertTrue((m0["reasoning_details"] as? JsonArray)?.isNotEmpty() == true)
}
@Test
fun reasoningDetailsTextIsCopiedIntoReasoningContent() {
// Только reasoning_details (без строки reasoning) — текст всё равно достаётся.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning_details":[{"type":"reasoning.text","text":"D"}]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("D", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun reasoningDetailsWithForeignTypeAreIgnored() {
// Чужой тип (reasoning.encrypted) игнорируется, берётся только reasoning.text.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning_details":[{"type":"reasoning.encrypted","text":"X"},{"type":"reasoning.text","text":"T"}]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("T", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun multipleReasoningDetailsAreJoinedInOrder() {
// Два reasoning.text склеиваются через "\n" в порядке массива.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning_details":[{"type":"reasoning.text","text":"a"},{"type":"reasoning.text","text":"b"}]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("a\nb", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun existingNonEmptyReasoningContentIsNotOverwritten() {
// Непустое нативное поле не перезаписываем текстом из рассуждений.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning":"R","reasoning_content":"нативное"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("нативное", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun assistantWithoutToolCallsIsUntouched() {
// Без tool_calls assistant-сообщение не трогаем (тело идентично).
val m = """{"role":"assistant","reasoning":"R","content":"hi"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
}
@Test
fun assistantWithEmptyToolCallsIsUntouched() {
// Пустой массив tool_calls = не трогаем (тело идентично).
val m = """{"role":"assistant","tool_calls":[],"reasoning":"R"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
}
@Test
fun userAndToolMessagesAreUntouched() {
// user/tool/system с теми же полями — не assistant, не трогаем.
val msgs = """
{"role":"user","tool_calls":[$TC],"reasoning":"R"},
{"role":"tool","tool_calls":[$TC],"reasoning":"R"},
{"role":"system","tool_calls":[$TC],"reasoning":"R"}
""".trimIndent().replace("\n", " ")
val out = apply(msgs, "reasoning_content", false)
val roles = outMsgs(out).map { str(it, "role") }
assertEquals(listOf("user", "tool", "system"), roles)
outMsgs(out).forEach { assertFalse(it.containsKey("reasoning_content")) }
}
@Test
fun emptyTextWithEmptyOkFillsEmptyString() {
// Рассуждений нет вовсе, но emptyOk=true — пишем пустую строку.
val m = """{"role":"assistant","tool_calls":[$TC]}"""
val out = apply(m, "reasoning_content", true)
val m0 = outMsgs(out)[0]
assertEquals("", str(m0, "reasoning_content"))
assertTrue(m0.containsKey("reasoning_content"))
}
@Test
fun emptyTextWithoutEmptyOkLeavesMessageUnchanged() {
// Рассуждений нет и emptyOk=false — сообщение не трогаем.
val m = """{"role":"assistant","tool_calls":[$TC]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
}
@Test
fun reasoningFieldNullLeavesBodyUntouched() {
// Поле не объявлено у провайдера (null) — тело не трогаем.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning":"R"}"""
val out = apply(m, null, false)
assertEquals("""{"messages":[$m]}""", out.toString())
}
@Test
fun messagesAbsentOrNotArrayLeavesBodyUntouched() {
// Нет `messages` — не трогаем.
val noMessages = json.parseToJsonElement("""{"model":"m1"}""").jsonObject
assertEquals(noMessages.toString(), applyReasoningField(noMessages, "reasoning_content", false).toString())
// `messages: null` (JsonNull) — тоже не трогаем.
val nullMessages = json.parseToJsonElement("""{"messages":null}""").jsonObject
val out = applyReasoningField(nullMessages, "reasoning_content", false)
assertEquals(nullMessages.toString(), out.toString())
}
@Test
fun messageOrderAndOtherFieldsArePreserved() {
// Меняется только 2-е (assistant с tool_calls), порядок и остальные поля на месте.
val msgs = """
{"role":"user","content":"hi"},
{"role":"assistant","tool_calls":[$TC],"reasoning":"R","content":""},
{"role":"tool","tool_call_id":"t1","content":"ok"},
{"role":"user","content":"again"}
""".trimIndent().replace("\n", " ")
val out = apply(msgs, "reasoning_content", false)
val list = outMsgs(out)
assertEquals(listOf("user", "assistant", "tool", "user"), list.map { str(it, "role") })
assertEquals("R", str(list[1], "reasoning_content"))
assertTrue((list[1]["tool_calls"] as? JsonArray)?.isNotEmpty() == true)
assertTrue(list[1].containsKey("reasoning"))
assertEquals("hi", str(list[0], "content"))
assertEquals("ok", str(list[2], "content"))
assertEquals("again", str(list[3], "content"))
}
}
@@ -0,0 +1,176 @@
package pw.binom.llmproxy
import io.ktor.utils.io.ByteChannel
import io.ktor.utils.io.close
import io.ktor.utils.io.readAvailable
import io.ktor.utils.io.writeFully
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlinx.coroutines.test.runTest
class StreamDoneContractTest {
/** Терминальный маркер, который должен оказаться (или не оказаться) в конце вывода. */
private val done = "data: [DONE]\n\n"
/** Прогнать сырой passthrough-поток через боевую обвязку и вернуть записанные байты как строку. */
private suspend fun runRaw(input: String): String {
val src = ByteChannel(autoFlush = true)
src.writeFully(input.encodeToByteArray())
src.close(null)
val out = ByteChannel(autoFlush = true)
out.streamRawWithDoneContract(src)
out.close(null)
return readAll(out)
}
/** Прогнать SSE через think-обвязку и вернуть записанные байты как строку. */
private suspend fun runThink(input: String, mode: String): String {
val src = ByteChannel(autoFlush = true)
src.writeFully(input.encodeToByteArray())
src.close(null)
val out = ByteChannel(autoFlush = true)
out.streamSseWithThinkTags(src, mode)
out.close(null)
return readAll(out)
}
private suspend fun readAll(ch: ByteChannel): String {
val sb = StringBuilder()
val buf = ByteArray(512)
while (true) {
val n = ch.readAvailable(buf)
if (n == -1) break
if (n > 0) sb.append(buf.decodeToString(0, n))
}
return sb.toString()
}
private fun countOccurrences(text: String, needle: String): Int {
var count = 0
var from = 0
while (true) {
val i = text.indexOf(needle, from)
if (i < 0) return count
count++
from = i + needle.length
}
}
@Test
fun textHelperSeesFinishReason() {
assertTrue(hasNonNullFinishReasonText("""{"choices":[{"finish_reason":"stop"}]}"""))
}
@Test
fun textHelperIgnoresNullFinishReason() {
assertFalse(hasNonNullFinishReasonText("""{"choices":[{"finish_reason":null}]}"""))
}
@Test
fun textHelperIgnoresEscapedOccurrenceInContent() {
// Экранированное вхождение внутри content — это не поле finish_reason.
assertFalse(hasNonNullFinishReasonText("""{"delta":{"content":"\"finish_reason\":\"stop\""}}"""))
}
@Test
fun detectorSeesFinishReasonSplitAcrossReads() {
// Маркер finish_reason разрезан границей чтения: детектор обязан его собрать.
val detector = StreamEndDetector()
detector.feed("""{"choices":[{"finish_rea""".encodeToByteArray())
detector.feed("""son":"stop"}]}""".encodeToByteArray())
assertTrue(detector.sawFinishReason)
assertTrue(detector.missingDoneMarker)
}
@Test
fun detectorSeesDoneSplitAcrossReads() {
// [DONE] разрезан границей чтения: детектор видит его и маркер не нужен.
val detector = StreamEndDetector()
detector.feed("data: [DO".encodeToByteArray())
detector.feed("NE]".encodeToByteArray())
assertTrue(detector.sawDone)
assertFalse(detector.missingDoneMarker)
}
@Test
fun detectorSeesDoneAndFinishReason() {
// В одном feed и finish_reason, и [DONE] — дописывать ничего не нужно.
val detector = StreamEndDetector()
detector.feed("""{"choices":[{"finish_reason":"stop"}]} data: [DONE]""".encodeToByteArray())
assertFalse(detector.missingDoneMarker)
}
@Test
fun detectorTruncatedStreamNeedsNoMarker() {
// Обрыв без finish_reason: маркер не дописываем (усечённый стрим честно остаётся усечённым).
val detector = StreamEndDetector()
detector.feed("""{"choices":[{"delta":{"content":"hi"}}]}""".encodeToByteArray())
assertFalse(detector.missingDoneMarker)
}
@Test
fun rawStreamIsByteExactAndAppendsDone() = runTest {
// Сырой поток с finish_reason и без [DONE]: байты входа не меняются, маркер дописывается ровно один.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"},\"finish_reason\":\"stop\"}]}\n\n"
val out = runRaw(input)
assertEquals(input + done, out, "выход должен быть вход + один маркер")
assertEquals(1, countOccurrences(out, done), "маркер должен быть ровно один")
}
@Test
fun rawStreamTruncatedPassesThroughUntouched() = runTest {
// Обрыв без finish_reason: ни одного [DONE], вход не меняется.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
val out = runRaw(input)
assertEquals(input, out, "усечённый поток должен уйти как есть")
assertEquals(0, countOccurrences(out, "[DONE]"))
}
@Test
fun rawStreamKeepsUpstreamDoneMarker() = runTest {
// Апстрим прислал свой [DONE]: второй маркер не дописывается.
val input = "data: {\"choices\":[{\"finish_reason\":\"stop\"}]}\n\n" + done
val out = runRaw(input)
assertEquals(1, countOccurrences(out, done), "маркер не должен дублироваться")
}
@Test
fun thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker() = runTest {
// think-обвязка: finish_reason есть, [DONE] нет — маркер дописывается ровно один.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"},\"finish_reason\":\"stop\"}]}\n\n"
val out = runThink(input, "split")
assertTrue(out.endsWith(done), "выход должен заканчиваться маркером: $out")
assertEquals(1, countOccurrences(out, done), "маркер должен быть ровно один")
}
@Test
fun thinkStreamTruncatedNeedsNoMarker() = runTest {
// think-обвязка: обрыв без finish_reason — маркер не дописываем.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
val out = runThink(input, "split")
assertEquals(0, countOccurrences(out, "[DONE]"), "при обрыве маркера быть не должно: $out")
}
@Test
fun thinkStreamKeepsSingleDone() = runTest {
// think-обвязка: finish_reason и собственный [DONE] — маркер остаётся ровно один.
val input = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"},\"finish_reason\":\"stop\"}]}\n\n" + done
val out = runThink(input, "split")
assertEquals(1, countOccurrences(out, done), "маркер не должен дублироваться")
}
}
@@ -0,0 +1,179 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class ThinkTagChunkTest {
// Литералы тегов OPEN/CLOSE из ThinkTagSplitter собираем из символьных
// кусочков, чтобы не записывать тег единой строкой в исходнике.
private val openTag = '<' + "think" + '>'
private val closeTag = '<' + "/think" + '>'
// Маркер рассуждения (5 букв) собираем из Unicode-кодов, без единой строки.
private val reasoning = listOf(0x0420, 0x0410, 0x0417, 0x0423, 0x041C)
.map { Char(it) }
.joinToString("")
private fun deltaOf(result: String, index: Int): JsonObject =
Json.parseToJsonElement(result).jsonObject
.get("choices")!!.jsonArray[index].jsonObject
.get("delta")!!.jsonObject
@Test
fun chunkSplitHoldsTailAcrossChunks() {
// Хвост открывающего тега, разрезанный на границе чанков, переживает
// два отдельных вызова через общий splitters (getOrPut по index).
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
// Чанк 1: текст + обрезанный префикс открывающего тега (не завершает тег).
val prefix = openTag.dropLast(1)
val obj1 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"текст$prefix"}}]}""",
).jsonObject
val r1 = transformThinkChunk(obj1, splitters, "split", true)
assertNotNull(r1)
val d1 = deltaOf(r1, 0)
assertEquals("текст", d1["content"]!!.jsonPrimitive.content)
assertFalse(d1.containsKey("reasoning_content"))
// Чанк 2 (тот же splitters): остаток открывающего тега + маркер.
val rest = openTag.last()
val obj2 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"$rest$reasoning"}}]}""",
).jsonObject
val r2 = transformThinkChunk(obj2, splitters, "split", true)
assertNotNull(r2)
val d2 = deltaOf(r2, 0)
assertFalse(d2.containsKey("content"))
assertEquals(reasoning, d2["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun chunkStripRemovesReasoning() {
// Режим strip: блок целиком (теги + рассуждения) вырезается,
// content склеивается в «AB», ключа reasoning_content нет.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val c = "A" + openTag + reasoning + closeTag + "B"
val obj = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"$c"}}]}""",
).jsonObject
val r = transformThinkChunk(obj, splitters, "strip", false)
assertNotNull(r)
val d = deltaOf(r, 0)
assertEquals("AB", d["content"]!!.jsonPrimitive.content)
assertFalse(d.containsKey("reasoning_content"))
}
@Test
fun chunkDropsEmptyContentFieldWhenAllGoesToReasoning() {
// Режим split: весь контент чанка — think-блок, поэтому content пуст
// и ключ убирается из delta, а рассуждения уходят в reasoning_content.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val c = openTag + reasoning + closeTag
val obj = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":"$c"}}]}""",
).jsonObject
val r = transformThinkChunk(obj, splitters, "split", true)
assertNotNull(r)
val d = deltaOf(r, 0)
assertFalse(d.containsKey("content"))
assertEquals(reasoning, d["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun chunkReturnsNullWhenNoChoices() {
// Нет ключа choices — чанк не трогаем, отдаём исходную строку (null).
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val obj = Json.parseToJsonElement("""{"model":"m"}""").jsonObject
assertNull(transformThinkChunk(obj, splitters, "split", true))
}
@Test
fun chunkReturnsNullWhenDeltaContentNotAString() {
// content = null (JsonNull) и content = массив частей — не строка,
// значит choice не трогаем, весь чанк не меняется -> null.
val s1 = mutableMapOf<Int, ThinkTagSplitter>()
val o1 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":null}}]}""",
).jsonObject
assertNull(transformThinkChunk(o1, s1, "split", true))
val s2 = mutableMapOf<Int, ThinkTagSplitter>()
val o2 = Json.parseToJsonElement(
"""{"choices":[{"index":0,"delta":{"content":[{"type":"text","text":"hi"}]}}]}""",
).jsonObject
assertNull(transformThinkChunk(o2, s2, "split", true))
}
@Test
fun chunkSeparateSplittersPerIndex() {
// Один объект splitters на все три вызова: по каждому index держим
// своего сплиттера, поэтому обрезанные хвосты по index не путаются.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val prefix = openTag.dropLast(1)
val rest = openTag.last()
// Вызов 1: два choice — index 0 «A»+префикс тега, index 1 «B»+тот же префикс.
val c1 = """{"choices":[{"index":0,"delta":{"content":"A$prefix"}},{"index":1,"delta":{"content":"B$prefix"}}]}"""
val r1 = transformThinkChunk(Json.parseToJsonElement(c1).jsonObject, splitters, "split", true)
assertNotNull(r1)
// Вызов 2: index 1 получает остаток тега + маркер → рассуждения у index 1.
val c2 = """{"choices":[{"index":1,"delta":{"content":"$rest${reasoning}1"}}]}"""
val r2 = transformThinkChunk(Json.parseToJsonElement(c2).jsonObject, splitters, "split", true)
assertNotNull(r2)
val d2 = deltaOf(r2, 0)
assertEquals(reasoning + "1", d2["reasoning_content"]!!.jsonPrimitive.content)
// Вызов 3: index 0 получает остаток тега + маркер → рассуждения у index 0.
val c3 = """{"choices":[{"index":0,"delta":{"content":"$rest${reasoning}0"}}]}"""
val r3 = transformThinkChunk(Json.parseToJsonElement(c3).jsonObject, splitters, "split", true)
assertNotNull(r3)
val d3 = deltaOf(r3, 0)
assertEquals(reasoning + "0", d3["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun chunkKeepsUntouchedChoiceAndFinishReason() {
// Первый choice целиком идёт через split; второй без ключа content
// (в нём только finish_reason) — остаётся нетронутым, как и был.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val c = "A" + openTag + reasoning + closeTag + "B"
val c1 = """{"choices":[{"index":0,"delta":{"content":"$c"}},{"index":1,"delta":{"finish_reason":"stop"}}]}"""
val r = transformThinkChunk(Json.parseToJsonElement(c1).jsonObject, splitters, "split", true)
assertNotNull(r)
val d0 = deltaOf(r, 0)
assertEquals("AB", d0["content"]!!.jsonPrimitive.content)
assertEquals(reasoning, d0["reasoning_content"]!!.jsonPrimitive.content)
val d1 = deltaOf(r, 1)
assertEquals("stop", d1["finish_reason"]!!.jsonPrimitive.content)
assertFalse(d1.containsKey("content"))
}
@Test
fun chunkIndexDefaultsToZeroWhenAbsent() {
// Ключ index отсутствует — оба раза используем сплиттер под индекс 0,
// поэтому хвост первого чанка доживал до рассуждений второго.
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val prefix = openTag.dropLast(1)
val rest = openTag.last()
val c1 = """{"choices":[{"delta":{"content":"A$prefix"}}]}"""
val r1 = transformThinkChunk(Json.parseToJsonElement(c1).jsonObject, splitters, "split", true)
assertNotNull(r1)
val c2 = """{"choices":[{"delta":{"content":"$rest$reasoning"}}]}"""
val r2 = transformThinkChunk(Json.parseToJsonElement(c2).jsonObject, splitters, "split", true)
assertNotNull(r2)
val d2 = deltaOf(r2, 0)
assertEquals(reasoning, d2["reasoning_content"]!!.jsonPrimitive.content)
}
}
@@ -0,0 +1,179 @@
package pw.binom.llmproxy
import io.ktor.utils.io.ByteChannel
import io.ktor.utils.io.close
import io.ktor.utils.io.readAvailable
import io.ktor.utils.io.writeFully
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class ThinkTagStreamTest {
/** Прогнать SSE-текст через боевую обвязку и вернуть то, что она записала. */
private suspend fun runStream(input: String, mode: String): String {
val src = ByteChannel(autoFlush = true)
src.writeFully(input.encodeToByteArray())
src.close(null)
val out = ByteChannel(autoFlush = true)
out.streamSseWithThinkTags(src, mode)
out.close(null)
val sb = StringBuilder()
val buf = ByteArray(512)
while (true) {
val n = out.readAvailable(buf)
if (n == -1) break
if (n > 0) sb.append(buf.decodeToString(0, n))
}
return sb.toString()
}
// Литералы тегов OPEN/CLOSE из ThinkTagSplitter собираем из символьных
// кусочков, чтобы не записывать тег единой строкой в исходнике.
private val openTag = '<' + "think" + '>'
private val closeTag = '<' + "/think" + '>'
// Маркер рассуждения (5 букв) собираем из Unicode-кодов, без единой строки.
private val reasoning = listOf(0x0420, 0x0410, 0x0417, 0x0423, 0x041C)
.map { Char(it) }
.joinToString("")
/** JSON-нагрузки из data-строк вывода (строки, начинающиеся с «data:»), без служебного [DONE]. */
private fun dataPayloads(sse: String): List<String> =
sse.lines()
.filter { it.startsWith("data:") && it.removePrefix("data:").trim() != "[DONE]" }
.map { it.removePrefix("data:").trim() }
private fun deltaOf(result: String, index: Int): JsonObject =
Json.parseToJsonElement(result).jsonObject
.get("choices")!!.jsonArray[index].jsonObject
.get("delta")!!.jsonObject
@Test
fun streamDoneAndServiceLinesPassThrough() = runTest {
// Служебные строки SSE и маркер конца потока должны дойти до клиента без изменений.
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n")
append("\n")
append("data: [DONE]\n")
append("\n")
append("event: ping\n")
append("\n")
append(": keep-alive\n")
append("\n")
}
val out = runStream(input, "split")
assertTrue(out.contains("data: [DONE]"), "маркер конца потока потерян: $out")
assertTrue(out.contains("event: ping"), "служебная строка event потеряна: $out")
assertTrue(out.contains(": keep-alive"), "строка-комментарий потеряна: $out")
assertTrue(out.contains("hi"), "текстовый чанк потерян: $out")
}
@Test
fun streamSplitsReasoningAcrossDataChunks() = runTest {
// Хвост открывающего тега разрезан на границе двух data-событий: первое
// отдаёт content «A» и хвост держит, второе завершает тег и блок.
val prefix = openTag.take(3)
val rest = openTag.drop(3)
val content2 = rest + reasoning + closeTag + "B"
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"A$prefix\"}}]}\n")
append("\n")
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"$content2\"}}]}\n")
append("\n")
}
val out = runStream(input, "split")
val payloads = dataPayloads(out)
val last = deltaOf(payloads.last(), 0)
assertEquals("B", last["content"]!!.jsonPrimitive.content)
assertEquals(reasoning, last["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun streamUnclosedBlockCarriesReasoningInSameChunk() = runTest {
// Незакрытый think-блок в конце потока: reasoning отдаётся сразу в том же
// чанке, где пришёл, финиш-чанка не появляется — удержанного хвоста нет.
val content = "A" + openTag + "МЫСЛИ"
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"$content\"}}]}\n")
append("\n")
}
val payloads = dataPayloads(runStream(input, "split"))
assertEquals(1, payloads.size, "финиш-чанка быть не должно: $payloads")
val delta = deltaOf(payloads.single(), 0)
assertEquals("A", delta["content"]!!.jsonPrimitive.content)
assertEquals("МЫСЛИ", delta["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun streamFlushesHeldTailOnFinish() = runTest {
// Поток обрывается на неполном префиксе открывающего тега: удержанный
// хвост не теряется и сбрасывается финиш-чанком в конце.
val tail = openTag.take(3)
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"A$tail\"}}]}\n")
append("\n")
}
val payloads = dataPayloads(runStream(input, "split"))
assertEquals(2, payloads.size, "финиш-чанк с удержанным хвостом потерян: $payloads")
val flushed = deltaOf(payloads.last(), 0)
assertEquals(tail, flushed["content"]!!.jsonPrimitive.content)
}
@Test
fun streamInvalidJsonPassesThroughVerbatim() = runTest {
// data-строка, которая не является JSON, доходит до клиента без изменений.
val input = buildString {
append("data: {это не json\n")
append("\n")
}
val out = runStream(input, "split")
assertTrue(out.contains("data: {это не json"), "битая строка не дошла как есть: $out")
}
@Test
fun streamChunkWithoutChoicesPassesThrough() = runTest {
// Чанк с пустыми choices не меняется и уходит клиенту исходной строкой.
val input = buildString {
append("data: {\"usage\":{\"total_tokens\":5},\"choices\":[]}\n")
append("\n")
}
val out = runStream(input, "split")
assertTrue(out.contains("\"total_tokens\":5"), "чанк с usage изменился: $out")
}
@Test
fun streamModeStripEmitsNoReasoning() = runTest {
// Режим strip: текст внутри think-тегов выбрасывается, поля
// reasoning_content нет, обычный текст остаётся.
val content = "A" + openTag + reasoning + closeTag + "B"
val input = buildString {
append("data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"$content\"}}]}\n")
append("\n")
}
val out = runStream(input, "strip")
val payloads = dataPayloads(out)
assertTrue(!out.contains("reasoning_content"), "в режиме strip не должно быть reasoning_content: $out")
assertEquals(1, payloads.size, "ожидался ровно один чанк: $payloads")
assertEquals("AB", deltaOf(payloads.single(), 0)["content"]!!.jsonPrimitive.content, "теги и рассуждение должны быть вырезаны: $payloads")
}
}
@@ -0,0 +1,174 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class ThinkTagTransformTest {
// Точные литералы тегов из ThinkTagSplitter (OPEN/CLOSE): собираем из
// отдельных символов, чтобы не писать тег единой строкой в исходнике.
private val openTag = '<' + "think" + '>'
private val closeTag = '<' + "/think" + '>'
private fun messageOf(result: String, index: Int): JsonObject =
Json.parseToJsonElement(result).jsonObject
.get("choices")!!.jsonArray[index].jsonObject
.get("message")!!.jsonObject
@Test
fun splitSingleBlockMovesReasoningAndKeepsOtherFields() {
// Проверяем: в режиме split содержимое блока уходит в reasoning_content,
// в content остаётся «AB», а служебные поля (model, usage, …) не тронуты.
val c = "A" + openTag + "РАЗУМ" + closeTag + "B"
val obj = Json.parseToJsonElement(
"""
{"id":"c1","object":"chat.completion","created":123,"model":"m1",
"usage":{"prompt_tokens":1,"completion_tokens":2,"total_tokens":3},
"system_fingerprint":"sf1",
"choices":[{"index":0,"message":{"role":"assistant","content":"$c"},"finish_reason":"stop"}]}
""".trimIndent(),
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val out = Json.parseToJsonElement(result).jsonObject
val msg = out.get("choices")!!.jsonArray[0].jsonObject.get("message")!!.jsonObject
assertEquals("AB", msg["content"]!!.jsonPrimitive.content)
assertEquals("РАЗУМ", msg["reasoning_content"]!!.jsonPrimitive.content)
assertEquals("m1", out["model"]!!.jsonPrimitive.content)
assertEquals(3, out["usage"]!!.jsonObject["total_tokens"]!!.jsonPrimitive.content.toInt())
}
@Test
fun splitTwoBlocksConcatenateInOrder() {
// Проверяем порядок конкатенации reasoning, когда два think-блока идут подряд.
val c = openTag + "ПЕРВЫЙ" + closeTag + "X" + openTag + "ВТОРОЙ" + closeTag
val obj = Json.parseToJsonElement(
"""{"choices":[{"index":0,"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("X", msg["content"]!!.jsonPrimitive.content)
assertEquals("ПЕРВЫЙВТОРОЙ", msg["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun splitAppendsToExistingReasoningContent() {
// Проверяем, что прежнее reasoning_content не теряется, а новое дописывается.
val c = openTag + "NEW" + closeTag
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c","reasoning_content":"OLD"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("OLDNEW", msg["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun splitUnclosedTagGoesToReasoning() {
// Проверяем, что незакрытый открывающий тег уводит хвост целиком в reasoning.
val c = "текст" + openTag + "мысли"
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("текст", msg["content"]!!.jsonPrimitive.content)
assertEquals("мысли", msg["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun stripRemovesTagsAndLeavesNoReasoningKey() {
// Проверяем: в режиме strip теги и рассуждения вырезаются,
// а ключ reasoning_content в message отсутствует вовсе.
val c = "A" + openTag + "РАЗУМ" + closeTag + "B"
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "strip")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("AB", msg["content"]!!.jsonPrimitive.content)
assertFalse(msg.containsKey("reasoning_content"))
}
@Test
fun missingChoicesReturnsNull() {
// Проверяем: без ключа choices функция не вносит изменений (null).
val obj = Json.parseToJsonElement(
"""{"id":"c1","model":"m1"}""",
).jsonObject
assertNull(transformThinkMessage(obj, "split"))
}
@Test
fun nonStringContentReturnsNull() {
// content = null (JsonNull) — менять нечего, функция возвращает null.
val nullContent = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":null}}]}""",
).jsonObject
assertNull(transformThinkMessage(nullContent, "split"))
// content = массив частей (JsonArray) — тоже не трогаем, null.
val arrayContent = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":[{"type":"text","text":"hi"}]}}]}""",
).jsonObject
assertNull(transformThinkMessage(arrayContent, "split"))
}
@Test
fun twoChoicesEachKeepOwnThinkBlock() {
// Проверяем, что каждый choice обрабатывается независимо: у каждого свой
// вырезанный content и свой reasoning_content.
val c0 = "А" + openTag + "А-раз" + closeTag + "Б"
val c1 = "В" + openTag + "В-раз" + closeTag + "Г"
val obj = Json.parseToJsonElement(
"""
{"choices":[
{"index":0,"message":{"content":"$c0"},"finish_reason":"stop"},
{"index":1,"message":{"content":"$c1"},"finish_reason":"stop"}]}
""".trimIndent(),
).jsonObject
val result = transformThinkMessage(obj, "split")
assertNotNull(result)
val out = Json.parseToJsonElement(result).jsonObject.get("choices")!!.jsonArray
val m0 = out[0].jsonObject.get("message")!!.jsonObject
val m1 = out[1].jsonObject.get("message")!!.jsonObject
assertEquals("АБ", m0["content"]!!.jsonPrimitive.content)
assertEquals("А-раз", m0["reasoning_content"]!!.jsonPrimitive.content)
assertEquals("ВГ", m1["content"]!!.jsonPrimitive.content)
assertEquals("В-раз", m1["reasoning_content"]!!.jsonPrimitive.content)
}
@Test
fun offModePassthroughKeepsTags() {
// Проверяем passthrough: в режиме off теги остаются в content,
// reasoning_content не добавляется, а результат — не null.
val c = "A" + openTag + "РАЗУМ" + closeTag + "B"
val obj = Json.parseToJsonElement(
"""{"choices":[{"message":{"content":"$c"}}]}""",
).jsonObject
val result = transformThinkMessage(obj, "off")
assertNotNull(result)
val msg = messageOf(result, 0)
assertEquals("A" + openTag + "РАЗУМ" + closeTag + "B", msg["content"]!!.jsonPrimitive.content)
assertFalse(msg.containsKey("reasoning_content"))
}
}