5 Commits

Author SHA1 Message Date
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
subochev c7e7c7346d feat: флажок think_tags — вырезание think-тегов из ответа (split/strip)
Build LLM Proxy / Build and push (release) Successful in 24s
Некоторые провайдеры (minimax) отдают рассуждения модели текстом внутри
content, обернув их тегами think. Новое опциональное поле think_tags
у providers[] и upstreams[] (приоритет: апстрим -> провайдер -> off,
толерантный разбор true/false, неизвестное = off).

- off: поведение не меняется (сырой байтовый passthrough в стриме);
- split: блоки think вырезаются из content, текст уходит в reasoning_content;
- strip: вырезаются и выбрасываются.

Стрим разбирается построчно: отдельный ThinkTagSplitter на каждый
choices[].index, хвост, который может оказаться началом разрезанного тега,
не отдаётся до разрешения, на конце потока — finish(). Non-stream: тот же
трансформ для message.content (и для результата rebuildFromChunks), уже
существующий reasoning_content дописывается, а не теряется.

Тесты: ThinkTagSplitterTest (7), разбор конфига и эффективный режим в
ConfigLogicTest (2) — всего 42 теста, 0 падений. CONFIG.md + README.
2026-09-11 16:24:54 +03:00
subochev dba52d22fd feat: вычисление сессии по истории (LCP) для providers[].session_header
- SHA-256 на commonMain + SessionRegistry (LCP по префикс-хэшам, LRU 1000/6ч, Mutex)
- providers[].session_header: прокси сам ставит/перезаписывает заголовок сессии
- цепочка хэшей стартует с первого user-сообщения (system не склеивает сессии)
- CONFIG.md + тесты (SHA-256, LCP, LRU/TTL, парсинг, стабильность префиксов)
2026-09-11 05:27:15 +03:00
subochev efdce75ee9 feat: логирование заголовков запроса и заголовков в апстрим; маскировка секретов; Content-Type прокси-овнер 2026-09-11 04:44:05 +03:00
16 changed files with 2024 additions and 12 deletions
+4
View File
@@ -9,6 +9,10 @@ jobs:
steps: steps:
- name: 'Checkout' - name: 'Checkout'
uses: https://github.com/actions/checkout@v4 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' - name: 'Build jar'
uses: https://git.binom.pw/subochev/devops/build-gradle@main uses: https://git.binom.pw/subochev/devops/build-gradle@main
with: with:
+88
View File
@@ -73,6 +73,7 @@ providers:
url: "https://routerai.ru/api/v1" url: "https://routerai.ru/api/v1"
key: "sk-..." # Bearer-ключ; можно подставлять из env key: "sk-..." # Bearer-ключ; можно подставлять из env
max_concurrency: 4 # опционально; лимит по умолчанию для апстримов max_concurrency: 4 # опционально; лимит по умолчанию для апстримов
session_header: x-opencode-session # опционально; прокси считает сессию из истории
patch: # уровень провайдера: ко всем его запросам patch: # уровень провайдера: ко всем его запросам
provider: provider:
allow_fallbacks: false allow_fallbacks: false
@@ -131,6 +132,92 @@ models:
`patch` опционален на **любом** уровне (`providers` / `upstreams` / `models`): `patch` опционален на **любом** уровне (`providers` / `upstreams` / `models`):
если ни одного нет — запрос проксируется как есть (исходное тело клиента). если ни одного нет — запрос проксируется как есть (исходное тело клиента).
### Сессия по истории (`session_header`)
`providers[].session_header` (опционально) — имя HTTP-заголовка, который прокси
**вычисляет сам** из истории сообщений и ставит в запрос к этому провайдеру.
Нужно для API, требующих стабильный идентификатор сессии (например,
`x-opencode-session`), когда клиент его не шлёт или шлёт не то.
```yaml
providers:
- id: some-provider
url: "https://.../v1"
session_header: x-opencode-session
```
Как считается id:
1. Берётся финальное тело запроса (после всех `patch`), из него — `messages`.
2. Цепочка **инкрементальных SHA-256 префикс-хэшей** начинается с первого
`user`-сообщения (ведущий `system`-промпт игнорируется: он обычно одинаков
у всех сессий клиента и как признак сессии бесполезен).
3. В реестре сессий ищется **наибольший общий префикс** (LCP) с уже виденной
историей. Нашли — используется id той сессии; не нашли — создаётся новая
(`id` = хэш всей истории на первом ходу).
4. Заголовок ставится **всегда** (клиентское значение перезаписывается).
Итог: пока история одной сессии растёт (дописываются assistant/user-сообщения),
id не меняется; разные диалоги получают разные id.
> **Ограничения.** Реестр живёт в памяти (LRU: 1000 сессий / 6 часов) — при
> рестарте прокси активные сессии получат новый id. Обрезка/суммаризация
> истории рвёт общий префикс → сессия распадётся на новую. Диалоги с
> одинаковым первым `user`-сообщением неразличимы (склеятся).
### Обработка think-тегов (`think_tags`)
Некоторые провайдеры (например, **minimax**) отдают рассуждения модели не в
отдельном поле `reasoning_content`, а прямо в `content`, обернув их тегами
`<think>…</think>`. Флажок `think_tags` (опционально) разрешает прокси разрезать
такой ответ и разложить его по полям.
Поле доступно на двух уровнях:
- `providers[].think_tags` — правило по умолчанию для всех апстримов провайдера;
- `upstreams[].think_tags` — необязательное переопределение на конкретной
апстрим-модели.
Значения (строки):
| Значение | Поведение |
|---|---|
| `off` | дефолт: ответ не меняется, теги остаются в `content` |
| `split` | блоки `<think>…</think>` вырезаются из `content`, их текст уходит в `reasoning_content` |
| `strip` | блоки вырезаются и выбрасываются — клиент рассуждений не видит |
Разбор значения толерантный: `true` ≡ `split`, `false` ≡ `off`; любое
неизвестное/пустое значение трактуется как `off` (прокси не падает).
Приоритет резолва — по убыванию специфичности: `upstreams[].think_tags` (самый
конкретный) → `providers[].think_tags` → `off`. То есть значение апстрима
перекрывает провайдерское.
```yaml
providers:
- id: minimax
url: "https://api.minimax.io/v1"
key: "${MINIMAX_API_KEY}"
think_tags: split # дефолт для всех апстримов провайдера
upstreams:
- id: minimax-m1
provider: minimax
model: MiniMax-M1 # наследует split от провайдера
- id: minimax-text-only
provider: minimax
model: MiniMax-Text-01
think_tags: strip # переопределение: рассуждения выбрасываем
```
Работает и в стриме, и в обычном (non-stream) ответе. Тег может прийти
разрезанным между чанками SSE — прокси держит хвост, который может оказаться
началом тега, и не отдаёт его клиенту до разрешения, поэтому огрызок тега не
утечёт. Незакрытый `<think>` в конце потока трактуется как «всё после него —
рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ
не трогаем. Без флажка (`off`) ответ идёт байт-в-байт как раньше.
### Пример сборки тела (многослойный `patch`) ### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -211,6 +298,7 @@ data class ProviderConf(
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, // заголовок-сессия, считается из истории
) )
@Serializable @Serializable
+5
View File
@@ -19,6 +19,11 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` (
умолчанию для его апстримов), и на апстриме (перекрывает провайдерский); без умолчанию для его апстримов), и на апстриме (перекрывает провайдерский); без
обоих — безлимит. обоих — безлимит.
Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или
апстрима (`off` по умолчанию, `split` — рассуждения из `<think>…</think>` уходят
в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в
non-stream.
| Переменная | Default | Описание | | Переменная | Default | Описание |
|---|---|---| |---|---|---|
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
+49
View File
@@ -0,0 +1,49 @@
# Тестирование 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` | Разбор конфига, приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `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`.
Правило: боевой код нельзя подгонять под тест; если тест не проходит, неверен тест.
## Известное ограничение
Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
+1
View File
@@ -31,6 +31,7 @@ kotlin {
val commonTest by getting { val commonTest by getting {
dependencies { dependencies {
implementation(libs.kotlin.test) implementation(libs.kotlin.test)
implementation(libs.kotlinx.coroutines.test)
} }
} }
val jvmMain by getting { 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" } yamlkt = { module = "net.mamoe.yamlkt:yamlkt", version.ref = "yamlkt" }
kotlinx-coroutines-core = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-core", version.ref = "coroutines" } 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-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" } kotlinx-io-core = { module = "org.jetbrains.kotlinx:kotlinx-io-core", version.ref = "kotlinxIo" }
logback-classic = { module = "ch.qos.logback:logback-classic", version.ref = "logback" } logback-classic = { module = "ch.qos.logback:logback-classic", version.ref = "logback" }
kotlin-test = { module = "org.jetbrains.kotlin:kotlin-test" } kotlin-test = { module = "org.jetbrains.kotlin:kotlin-test" }
+319 -12
View File
@@ -3,7 +3,9 @@ package pw.binom.llmproxy
import io.ktor.client.HttpClient import io.ktor.client.HttpClient
import io.ktor.client.call.body import io.ktor.client.call.body
import io.ktor.client.request.headers import io.ktor.client.request.headers
import io.ktor.utils.io.LineEnding
import io.ktor.utils.io.readAvailable import io.ktor.utils.io.readAvailable
import io.ktor.utils.io.readLine
import io.ktor.http.Headers import io.ktor.http.Headers
import io.ktor.client.request.preparePost import io.ktor.client.request.preparePost
import io.ktor.client.request.setBody import io.ktor.client.request.setBody
@@ -14,7 +16,9 @@ import io.ktor.server.application.Application
import io.ktor.server.application.ApplicationCall import io.ktor.server.application.ApplicationCall
import io.ktor.server.application.call import io.ktor.server.application.call
import io.ktor.server.application.install import io.ktor.server.application.install
import io.ktor.server.request.httpMethod
import io.ktor.server.request.receiveText import io.ktor.server.request.receiveText
import io.ktor.server.request.uri
import io.ktor.server.response.respondBytesWriter import io.ktor.server.response.respondBytesWriter
import io.ktor.server.response.respondText import io.ktor.server.response.respondText
import io.ktor.server.routing.get import io.ktor.server.routing.get
@@ -95,8 +99,9 @@ fun main() {
} }
val http = createHttpClient() val http = createHttpClient()
val sessions = SessionRegistry()
startServer(config.server.host, config.server.port) { startServer(config.server.host, config.server.port) {
proxyModule(config, providersById, upstreamsById, active, http) proxyModule(config, providersById, upstreamsById, active, sessions, http)
} }
} }
@@ -107,11 +112,12 @@ fun Application.proxyModule(
providersById: Map<String, ProviderConf>, providersById: Map<String, ProviderConf>,
upstreamsById: Map<String, UpstreamConf>, upstreamsById: Map<String, UpstreamConf>,
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
) { ) {
routing { routing {
post("/v1/chat/completions") { post("/v1/chat/completions") {
handleChat(call, config, providersById, upstreamsById, active, http) handleChat(call, config, providersById, upstreamsById, active, sessions, http)
} }
get("/v1/models") { get("/v1/models") {
handleModels(call, config) handleModels(call, config)
@@ -125,8 +131,10 @@ private suspend fun handleChat(
providersById: Map<String, ProviderConf>, providersById: Map<String, ProviderConf>,
upstreamsById: Map<String, UpstreamConf>, upstreamsById: Map<String, UpstreamConf>,
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
) { ) {
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val raw = call.receiveText() val raw = call.receiveText()
if (raw.isBlank()) { if (raw.isBlank()) {
call.respondText(errorJson("empty body"), ContentType.Application.Json, HttpStatusCode.BadRequest) call.respondText(errorJson("empty body"), ContentType.Application.Json, HttpStatusCode.BadRequest)
@@ -183,13 +191,27 @@ private suspend fun handleChat(
val url = provider.url.trimEnd('/') + "/chat/completions" val url = provider.url.trimEnd('/') + "/chat/completions"
val providerKey = resolveEnv(provider.key) val providerKey = resolveEnv(provider.key)
val sessionHeader = provider.session_header
// Клиентский одноимённый заголовок не пробрасываем: при метке
// значение x-opencode-session всегда вычисляем сами (LCP по истории).
val forwardedHeaders = headersToForward(call.request.headers)
.filterKeys { sessionHeader == null || !it.equals(sessionHeader, ignoreCase = true) }
val sessionId = sessionHeader?.let { sessions.resolve(sessionPrefixHashes(forwarded)) }
val outgoingHeaders = forwardedHeaders.toMutableMap().apply {
this["Content-Type"] = listOf("application/json")
if (providerKey.isNotEmpty()) this["Authorization"] = listOf("Bearer $providerKey")
if (sessionHeader != null && sessionId != null) this[sessionHeader] = listOf(sessionId)
}
log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} session=${sessionId ?: "-"} → $url headers: ${formatHeadersForLog(outgoingHeaders)}"
}
var failover = false var failover = false
var responded = false var responded = false
var upstreamStatus = 0 var upstreamStatus = 0
http.preparePost(url) { http.preparePost(url) {
headers { headers {
headersToForward(call.request.headers).forEach { (name, values) -> forwardedHeaders.forEach { (name, values) ->
appendAll(name, values) appendAll(name, values)
} }
// Авторизация — всегда наша (ключ провайдера из конфига); // Авторизация — всегда наша (ключ провайдера из конфига);
@@ -199,6 +221,9 @@ private suspend fun handleChat(
set("Authorization", "Bearer $providerKey") set("Authorization", "Bearer $providerKey")
} }
set("Content-Type", "application/json") set("Content-Type", "application/json")
if (sessionHeader != null && sessionId != null) {
set(sessionHeader, sessionId)
}
} }
setBody(forwarded.toString()) setBody(forwarded.toString())
}.execute { resp -> }.execute { resp ->
@@ -213,16 +238,14 @@ private suspend fun handleChat(
if (clientWantsStream) { if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream" val ct = resp.headers["Content-Type"] ?: "text/event-stream"
val status = HttpStatusCode.fromValue(upstreamStatus) val status = HttpStatusCode.fromValue(upstreamStatus)
val thinkMode = effectiveThinkTags(up, provider)
call.respondBytesWriter(ContentType.parse(ct), status) { call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>() val ch = resp.body<ByteReadChannel>()
val buf = ByteArray(8192) if (thinkMode == "off") {
while (true) { // Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
val n = ch.readAvailable(buf) streamRawWithDoneContract(ch)
if (n == -1) break } else {
if (n > 0) { streamSseWithThinkTags(ch, thinkMode)
writeFully(buf, 0, n)
flush()
}
} }
} }
log.info { log.info {
@@ -232,7 +255,12 @@ private suspend fun handleChat(
} else { } else {
val ct = resp.headers["Content-Type"] ?: "application/json" val ct = resp.headers["Content-Type"] ?: "application/json"
val full = resp.body<String>() val full = resp.body<String>()
val out = if (ct.contains("text/event-stream")) rebuildFromChunks(full) else full val thinkMode = effectiveThinkTags(up, provider)
val rebuilt = if (ct.contains("text/event-stream")) rebuildFromChunks(full) else full
val out = if (thinkMode != "off") (
runCatching { transformThinkMessage(json.parseToJsonElement(rebuilt).jsonObject, thinkMode) }
.getOrNull() ?: rebuilt
) else rebuilt
val outCt = runCatching { val outCt = runCatching {
if (ct.contains("text/event-stream") && if (ct.contains("text/event-stream") &&
Json.parseToJsonElement(out).jsonObject["error"] != null Json.parseToJsonElement(out).jsonObject["error"] != null
@@ -287,6 +315,9 @@ private val SKIP_HEADER_NAMES = setOf(
// Authorization управляется прокси явно (ключ провайдера), клиентский // Authorization управляется прокси явно (ключ провайдера), клиентский
// не пересылается // не пересылается
"authorization", "authorization",
// Content-Type всегда наш (application/json: тело мержится как JSON),
// клиентский не пересылаем, чтобы не ушло двух заголовков
"content-type",
) )
/** /**
@@ -302,6 +333,29 @@ internal fun headersToForward(request: Headers): Map<String, List<String>> =
.filter { (name, _) -> name.lowercase() !in SKIP_HEADER_NAMES } .filter { (name, _) -> name.lowercase() !in SKIP_HEADER_NAMES }
.associate { (name, values) -> name to values } .associate { (name, values) -> name to values }
/** Заголовки, значения которых маскируются в логах (секреты клиента). */
private val SENSITIVE_HEADER_NAMES = setOf(
"authorization",
"proxy-authorization",
"x-api-key",
"api-key",
"cookie",
"set-cookie",
)
/**
* Заголовки в виде строки для лога (`name=v1|v2, ...`). Значения чувствительных
* имён ([SENSITIVE_HEADER_NAMES]) маскируются `***`, чтобы не светить секреты.
*/
internal fun formatHeadersForLog(headers: Map<String, List<String>>): String =
headers.entries.joinToString(", ") { (name, values) ->
val shown = if (name.lowercase() in SENSITIVE_HEADER_NAMES) values.map { "***" } else values
"$name=${shown.joinToString("|")}"
}
internal fun formatHeadersForLog(headers: Headers): String =
formatHeadersForLog(headers.entries().associate { (name, values) -> name to values })
/** /**
* Сборка тела запроса: подмена `model` на реальное имя апстрима + глубокий * Сборка тела запроса: подмена `model` на реальное имя апстрима + глубокий
* послойный мерж `patch` в порядке provider → upstream → model. * послойный мерж `patch` в порядке provider → upstream → model.
@@ -367,6 +421,20 @@ internal fun tryClaim(up: UpstreamConf, active: Map<String, UpstreamCounter>): B
internal fun effectiveConcurrencyLimit(up: UpstreamConf, provider: ProviderConf?): Int = internal fun effectiveConcurrencyLimit(up: UpstreamConf, provider: ProviderConf?): Int =
up.max_concurrency ?: provider?.max_concurrency ?: Int.MAX_VALUE up.max_concurrency ?: provider?.max_concurrency ?: Int.MAX_VALUE
/**
* Эффективные think_tags апстрима: значение у апстрима, если задано; иначе у
* провайдера; иначе "off". Значения "true" трактуются как "split", "false" и
* любое неизвестное/пустое — как "off".
*/
internal fun effectiveThinkTags(up: UpstreamConf, provider: ProviderConf?): String {
val raw = up.think_tags ?: provider?.think_tags ?: "off"
return when (raw) {
"split", "strip" -> raw
"true" -> "split"
else -> "off"
}
}
/** Освободить слот апстрима (в finally по завершении проксирования). */ /** Освободить слот апстрима (в finally по завершении проксирования). */
internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) { internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) {
active.getValue(up.id).release() active.getValue(up.id).release()
@@ -422,6 +490,7 @@ internal fun pickFreeUpstream(
pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) } pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) }
private suspend fun handleModels(call: ApplicationCall, config: Config) { private suspend fun handleModels(call: ApplicationCall, config: Config) {
log.info { "[llm-proxy] models ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val created = TimeSource.Monotonic.markNow().elapsedNow().inWholeSeconds val created = TimeSource.Monotonic.markNow().elapsedNow().inWholeSeconds
val data = config.models.map { m -> val data = config.models.map { m ->
JsonObject( JsonObject(
@@ -530,6 +599,238 @@ internal fun rebuildFromChunks(sse: String): String {
return JsonObject(root).toString() return JsonObject(root).toString()
} }
/**
* Пересобрать полный chat.completion (обычный JSON или результат
* [rebuildFromChunks]): по каждому choice взять `message.content` — только если
* это JSON-строка (массив частей не трогаем) — и прогнать целиком через
* [ThinkTagSplitter] (feed + finish). Остаток возвращается в `message.content`,
* вырезанное — в `message.reasoning_content` (режим split; при strip не
* добавляем). Если `reasoning_content` уже был непустой строкой — новое
* дописываем в конец существующего, не теряя прежнее. Ни один choice не
* изменился → null (отдать исходную строку как есть).
*/
internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String? {
val choices = obj["choices"]?.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 contentStr = (message["content"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl
val splitter = ThinkTagSplitter(thinkMode)
val (outContent, outReasoning) = splitter.feed(contentStr)
val (tailContent, tailReasoning) = splitter.finish()
val reasoning = outReasoning + tailReasoning
changed = true
val newMessage = message.toMutableMap().apply {
this["content"] = JsonPrimitive(outContent + tailContent)
if (addReasoning && reasoning.isNotEmpty()) {
val existing = (this["reasoning_content"] as? JsonPrimitive)?.takeIf { it.isString }?.content
this["reasoning_content"] = JsonPrimitive((existing ?: "") + reasoning)
}
}
JsonObject(choice.toMutableMap().apply { this["message"] = JsonObject(newMessage) })
}
if (!changed) return null
return JsonObject(obj.toMutableMap().apply { this["choices"] = JsonArray(newChoices) }).toString()
}
/** Терминальный маркер 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()
}
}
/**
* Построчный разбор SSE-стрима с рассечением think-тегов. Строки, не начинающиеся
* с `data:`, и `data: [DONE]` уходят клиенту без изменений (с `\n`). Прочие
* `data:`-строки парсятся и прогоняются через [transformThinkChunk]; результат
* записывается как `data: <json>\n\n` (событие-граница SSE), а при ошибке парса —
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
*/
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()
}
if (thinkMode == "split" || thinkMode == "strip") {
splitters.forEach { (idx, sp) ->
val (tail, reasoning) = sp.finish()
val hasContent = tail.isNotEmpty()
val hasReasoning = addReasoning && reasoning.isNotEmpty()
if (!hasContent && !hasReasoning) return@forEach
val delta = mutableMapOf<String, JsonElement>()
if (hasContent) delta["content"] = JsonPrimitive(tail)
if (hasReasoning) delta["reasoning_content"] = JsonPrimitive(reasoning)
val chunk = JsonObject(
mutableMapOf(
"choices" to JsonArray(
listOf(
JsonObject(
mutableMapOf(
"index" to JsonPrimitive(idx),
"delta" to JsonObject(delta),
),
),
),
),
),
)
emitUtf8("data: $chunk\n\n")
flush()
}
}
// После сброса хвостов дописываем маркер, только если апстрим его не прислал,
// а `finish_reason` в потоке был (при обрыве без него маркер НЕ дописываем).
if (sawFinishReason && !sawDone) {
emitUtf8(SSE_DONE_MARKER)
flush()
}
}
/**
* Пересобрать SSE-чанк: по каждому choice (ключ `index`, дефолт 0) взять
* `delta.content` (только если это JSON-строка; массив частей не трогаем) и
* прогнать через [ThinkTagSplitter] для этого index. Остаток возвращается в
* `delta.content` (поле убирается, если пустое); вырезанное — в
* `delta.reasoning_content` (только режим split, при strip не добавляем).
* Чанк без `choices` или без строкового `delta.content` не меняется —
* возвращается null (отдать исходную строку как есть).
*/
internal fun transformThinkChunk(
obj: JsonObject,
splitters: MutableMap<Int, ThinkTagSplitter>,
thinkMode: String,
addReasoning: Boolean,
): String? {
val choices = obj["choices"]?.jsonArray ?: return null
var changed = false
val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject
val delta = choice["delta"]?.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 splitter = splitters.getOrPut(idx) { ThinkTagSplitter(thinkMode) }
val (newContent, reasoning) = splitter.feed(contentStr)
changed = true
val newDelta = delta.toMutableMap()
if (newContent.isEmpty()) newDelta.remove("content") else newDelta["content"] = JsonPrimitive(newContent)
if (addReasoning && reasoning.isNotEmpty()) newDelta["reasoning_content"] = JsonPrimitive(reasoning)
JsonObject(choice.toMutableMap().apply { this["delta"] = JsonObject(newDelta) })
}
if (!changed) return null
return JsonObject(obj.toMutableMap().apply { this["choices"] = JsonArray(newChoices) }).toString()
}
/** Записать строку как UTF-8 байты (KMP-безопасно, без java.io). */
private suspend fun ByteWriteChannel.emitUtf8(text: String) {
val bytes = text.encodeToByteArray()
writeFully(bytes, 0, bytes.size)
}
@Serializable @Serializable
data class ProviderConf( data class ProviderConf(
val id: String, val id: String,
@@ -537,6 +838,8 @@ data class ProviderConf(
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 think_tags: String? = null,
) )
data class UpstreamConf( data class UpstreamConf(
@@ -545,6 +848,7 @@ data class UpstreamConf(
val model: String, val model: String,
val max_concurrency: Int? = null, val max_concurrency: Int? = null,
val patch: JsonObject? = null, val patch: JsonObject? = null,
val think_tags: String? = null,
) )
data class ModelConf( data class ModelConf(
@@ -598,6 +902,8 @@ internal fun parseConfig(root: YamlElement): Config {
key = m.strOrNull("key") ?: "", key = m.strOrNull("key") ?: "",
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
session_header = m.strOrNull("session_header"),
think_tags = m.strOrNull("think_tags"),
) )
} }
@@ -609,6 +915,7 @@ internal fun parseConfig(root: YamlElement): Config {
model = m.str("model"), model = m.str("model"),
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
think_tags = m.strOrNull("think_tags"),
) )
} }
@@ -0,0 +1,206 @@
package pw.binom.llmproxy
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.datetime.Clock
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonElement
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive
/**
* SHA-256, чистая реализация на commonMain (без внешних зависимостей —
* KMP-зеркало их не отдаёт). Нужен для стабильного хэша префиксов истории.
*/
internal object Sha256 {
private val K = intArrayOf(
0x428a2f98, 0x71374491, 0xb5c0fbcf.toInt(), 0xe9b5dba5.toInt(), 0x3956c25b, 0x59f111f1, 0x923f82a4.toInt(), 0xab1c5ed5.toInt(),
0xd807aa98.toInt(), 0x12835b01, 0x243185be, 0x550c7dc3, 0x72be5d74, 0x80deb1fe.toInt(), 0x9bdc06a7.toInt(), 0xc19bf174.toInt(),
0xe49b69c1.toInt(), 0xefbe4786.toInt(), 0x0fc19dc6, 0x240ca1cc, 0x2de92c6f, 0x4a7484aa, 0x5cb0a9dc, 0x76f988da,
0x983e5152.toInt(), 0xa831c66d.toInt(), 0xb00327c8.toInt(), 0xbf597fc7.toInt(), 0xc6e00bf3.toInt(), 0xd5a79147.toInt(), 0x06ca6351, 0x14292967,
0x27b70a85, 0x2e1b2138, 0x4d2c6dfc, 0x53380d13, 0x650a7354, 0x766a0abb, 0x81c2c92e.toInt(), 0x92722c85.toInt(),
0xa2bfe8a1.toInt(), 0xa81a664b.toInt(), 0xc24b8b70.toInt(), 0xc76c51a3.toInt(), 0xd192e819.toInt(), 0xd6990624.toInt(), 0xf40e3585.toInt(), 0x106aa070,
0x19a4c116, 0x1e376c08, 0x2748774c, 0x34b0bcb5, 0x391c0cb3, 0x4ed8aa4a, 0x5b9cca4f, 0x682e6ff3,
0x748f82ee, 0x78a5636f, 0x84c87814.toInt(), 0x8cc70208.toInt(), 0x90befffa.toInt(), 0xa4506ceb.toInt(), 0xbef9a3f7.toInt(), 0xc67178f2.toInt(),
)
private fun rotr(x: Int, n: Int): Int = (x ushr n) or (x shl (32 - n))
private fun pad(input: ByteArray): ByteArray {
val total = ((input.size + 9 + 63) / 64) * 64
val out = ByteArray(total)
input.copyInto(out)
out[input.size] = 0x80.toByte()
val bits = input.size.toLong() * 8
for (k in 0 until 8) {
out[total - 1 - k] = (bits ushr (8 * k)).toByte()
}
return out
}
fun hash(input: ByteArray): ByteArray {
val msg = pad(input)
val h = intArrayOf(
0x6a09e667, 0xbb67ae85.toInt(), 0x3c6ef372, 0xa54ff53a.toInt(),
0x510e527f, 0x9b05688c.toInt(), 0x1f83d9ab, 0x5be0cd19,
)
val w = IntArray(64)
var i = 0
while (i < msg.size) {
for (j in 0 until 16) {
w[j] = ((msg[i + j * 4].toInt() and 0xff) shl 24) or
((msg[i + j * 4 + 1].toInt() and 0xff) shl 16) or
((msg[i + j * 4 + 2].toInt() and 0xff) shl 8) or
(msg[i + j * 4 + 3].toInt() and 0xff)
}
for (j in 16 until 64) {
val s0 = rotr(w[j - 15], 7) xor rotr(w[j - 15], 18) xor (w[j - 15] ushr 3)
val s1 = rotr(w[j - 2], 17) xor rotr(w[j - 2], 19) xor (w[j - 2] ushr 10)
w[j] = w[j - 16] + s0 + w[j - 7] + s1
}
var a = h[0]; var b = h[1]; var c = h[2]; var d = h[3]
var e = h[4]; var f = h[5]; var g = h[6]; var hh = h[7]
for (j in 0 until 64) {
val s1 = rotr(e, 6) xor rotr(e, 11) xor rotr(e, 25)
val ch = (e and f) xor (e.inv() and g)
val t1 = hh + s1 + ch + K[j] + w[j]
val s0 = rotr(a, 2) xor rotr(a, 13) xor rotr(a, 22)
val maj = (a and b) xor (a and c) xor (b and c)
val t2 = s0 + maj
hh = g; g = f; f = e; e = d + t1; d = c; c = b; b = a; a = t1 + t2
}
h[0] += a; h[1] += b; h[2] += c; h[3] += d
h[4] += e; h[5] += f; h[6] += g; h[7] += hh
i += 64
}
val out = ByteArray(32)
for (j in 0 until 8) {
out[j * 4] = (h[j] ushr 24).toByte()
out[j * 4 + 1] = (h[j] ushr 16).toByte()
out[j * 4 + 2] = (h[j] ushr 8).toByte()
out[j * 4 + 3] = h[j].toByte()
}
return out
}
}
private const val HEX = "0123456789abcdef"
/** SHA-256 строки (UTF-8) в нижнем hex. */
internal fun sha256Hex(text: String): String {
val bytes = Sha256.hash(text.encodeToByteArray())
val sb = StringBuilder(bytes.size * 2)
for (b in bytes) {
val v = b.toInt() and 0xff
sb.append(HEX[v ushr 4]).append(HEX[v and 0xf])
}
return sb.toString()
}
/**
* Инкрементальные префикс-хэши истории: H_i = sha256(H_{i-1} + "\u0000" + msg_i).
*
* Цепочка начинается с первого `user`-сообщения: system-промпт обычно
* идентичен у всех сессий одного клиента и как признак сессии бесполезен
* (иначе LCP склеивает все сессии на общем префиксе `[system]`).
* Если `messages` нет/пусто — fallback на хэш всего тела.
*/
internal fun sessionPrefixHashes(body: JsonObject): List<String> {
val messages = body["messages"] as? JsonArray
if (messages == null || messages.isEmpty()) {
return listOf(sha256Hex(body.toString()))
}
fun roleAt(i: Int): String? =
((messages[i] as? JsonObject)?.get("role") as? JsonPrimitive)?.content
var start = messages.indices.firstOrNull { roleAt(it) == "user" }
?: messages.indices.firstOrNull { roleAt(it) != "system" }
?: 0
val res = ArrayList<String>(messages.size - start)
var prev = ""
for (i in start until messages.size) {
prev = sha256Hex(prev + "\u0000" + messages[i].toString())
res.add(prev)
}
return res
}
/**
* Реестр сессий: id сессии определяется наибольшим общим префиксом (LCP)
* присланной истории. Для префиксов, которые уже встречались, возвращается
* id исходной сессии; иначе создаётся новая (id = хэш всей истории).
*
* Таблица ограничена по размеру (LRU) и времени жизни (TTL). Доступ под
* Mutex: параллельные запросы одной сессии не должны гонять состояние.
*/
class SessionRegistry(
private val maxSessions: Int = 1000,
private val ttlMillis: Long = 6 * 60 * 60 * 1000L,
) {
private class Entry(val prefixes: MutableSet<String>) {
var lastAccess: Long = 0
}
private val mutex = Mutex()
private val prefixToSession = mutableMapOf<String, String>()
/** В порядке доступа: голова — самая давняя, хвост — свежая. */
private val sessions = LinkedHashMap<String, Entry>()
suspend fun resolve(prefixHashes: List<String>): String? =
mutex.withLock { resolveLocked(prefixHashes, Clock.System.now().toEpochMilliseconds()) }
/** Синхронная (без блокировки) версия — для тестов и вызовов под mutex. */
internal fun resolveLocked(prefixHashes: List<String>, now: Long): String? {
if (prefixHashes.isEmpty()) return null
evictExpired(now)
var found: String? = null
for (i in prefixHashes.indices.reversed()) {
val s = prefixToSession[prefixHashes[i]]
if (s != null) {
found = s
break
}
}
val id = found ?: prefixHashes.last()
val entry = sessions.remove(id) ?: Entry(mutableSetOf())
for (h in prefixHashes) {
val prev = prefixToSession.put(h, id)
if (prev != null && prev != id) {
sessions[prev]?.prefixes?.remove(h)
}
entry.prefixes.add(h)
}
entry.lastAccess = now
sessions[id] = entry
evictOverflow()
return id
}
private fun drop(sessionId: String, entry: Entry) {
entry.prefixes.forEach { h -> if (prefixToSession[h] == sessionId) prefixToSession.remove(h) }
}
private fun evictExpired(now: Long) {
val it = sessions.entries.iterator()
while (it.hasNext()) {
val e = it.next()
if (now - e.value.lastAccess <= ttlMillis) break
drop(e.key, e.value)
it.remove()
}
}
private fun evictOverflow() {
val it = sessions.entries.iterator()
while (sessions.size > maxSessions && it.hasNext()) {
val e = it.next()
drop(e.key, e.value)
it.remove()
}
}
}
@@ -0,0 +1,96 @@
package pw.binom.llmproxy
/**
* Автомат рассечения кусков стрима по think-тегам (строго lowercase). Чистая логика: без Ktor, без IO, без побочных
* эффектов — по одному экземпляру на choice/index.
*
* - `off` — passthrough: всё, включая теги, уходит в content;
* - `split` — текст внутри think-блока → reasoning,
* остальное → content; несколько блоков в одном ответе обрабатываются все;
* - `strip` — как `split`, но рассуждения выбрасываются (reasoning всегда "").
*
* Тег может прийти разрезанным между кусками: хвост, который является
* префиксом ожидаемого тега: вне блока — <think>, внутри блока — </think>.
* Такой хвост не отдаём, держим до следующего `feed`. Если
* кусок показал, что хвост не тег, — отдаём его как обычный текст.
* Незакрытый think-блок в конце: всё после него (и накопленный хвост) — reasoning,
* сбрасывается в `finish`. Незнакомый режим ведёт себя как `off`.
*/
class ThinkTagSplitter(private val mode: String) {
private val active: Boolean = mode == "split" || mode == "strip"
private var inThink = false
private var pending = ""
fun feed(text: String): Pair<String, String> {
if (!active) return text to ""
val step = process(pending + text)
pending = step.pending
inThink = step.inThink
return step.content to if (mode == "strip") "" else step.reasoning
}
fun finish(): Pair<String, String> {
if (!active) return "" to ""
val tail = pending
val wasInThink = inThink
pending = ""
inThink = false
if (wasInThink) {
return "" to if (mode == "strip") "" else tail
}
return tail to ""
}
private class Step(
val content: String,
val reasoning: String,
val pending: String,
val inThink: Boolean,
)
private fun process(s: String): Step {
val content = StringBuilder()
val reasoning = StringBuilder()
var i = 0
var think = inThink
var hold = ""
while (i < s.length) {
val tag = if (think) CLOSE else OPEN
val idx = s.indexOf(tag, i)
if (idx < 0) {
val keep = holdableSuffix(s.substring(i), tag)
if (think) {
reasoning.append(s, i, s.length - keep)
} else {
content.append(s, i, s.length - keep)
}
hold = s.substring(s.length - keep)
break
}
if (think) {
reasoning.append(s, i, idx)
} else {
content.append(s, i, idx)
}
think = !think
i = idx + tag.length
}
return Step(content.toString(), reasoning.toString(), hold, think)
}
/** Максимальный суффикс хвоста, который является префиксом тега (0..tag.length-1). */
private fun holdableSuffix(tail: String, tag: String): Int {
var k = minOf(tag.length - 1, tail.length)
while (k > 0 && !tail.endsWith(tag.substring(0, k))) {
k--
}
return k
}
private companion object {
const val OPEN = "<think>"
const val CLOSE = "</think>"
}
}
@@ -289,6 +289,7 @@ class ConfigLogicTest {
"Proxy-Connection" to listOf("keep-alive"), "Proxy-Connection" to listOf("keep-alive"),
"Upgrade" to listOf("h2c"), "Upgrade" to listOf("h2c"),
"Authorization" to listOf("Bearer client-secret"), "Authorization" to listOf("Bearer client-secret"),
"Content-Type" to listOf("application/x-www-form-urlencoded"),
"Accept" to listOf("*/*"), "Accept" to listOf("*/*"),
) )
val out = headersToForward(req) val out = headersToForward(req)
@@ -297,6 +298,7 @@ class ConfigLogicTest {
assertEquals(listOf("*/*"), out["Accept"]) assertEquals(listOf("*/*"), out["Accept"])
assertEquals(3, out.size) assertEquals(3, out.size)
assertEquals(null, out["Authorization"]) assertEquals(null, out["Authorization"])
assertEquals(null, out["Content-Type"])
} }
@Test @Test
@@ -312,6 +314,32 @@ class ConfigLogicTest {
assertEquals(listOf("s"), out["x-opencode-session"]) assertEquals(listOf("s"), out["x-opencode-session"])
} }
@Test
fun formatHeadersForLogMasksSecretsAndKeepsOthers() {
val line = formatHeadersForLog(
headersOf(
"X-Opencode-Session" to listOf("abc-123"),
"Authorization" to listOf("Bearer super-secret"),
"x-api-key" to listOf("key-1"),
"Cookie" to listOf("session=deadbeef"),
),
)
assertTrue(line.contains("X-Opencode-Session=abc-123"))
assertTrue(line.contains("Authorization=***"))
assertTrue(line.contains("x-api-key=***"))
assertTrue(line.contains("Cookie=***"))
assertFalse(line.contains("super-secret"))
assertFalse(line.contains("deadbeef"))
}
@Test
fun formatHeadersForLogJoinsMultipleValues() {
val line = formatHeadersForLog(
headersOf("X-Custom" to listOf("a", "b")),
)
assertEquals("X-Custom=a|b", line)
}
@Test @Test
fun rebuildFromChunksPreservesAllUpstreamFields() { fun rebuildFromChunksPreservesAllUpstreamFields() {
val sse = """ val sse = """
@@ -354,4 +382,168 @@ class ConfigLogicTest {
val err = out["error"]?.jsonObject ?: error("error block missing") val err = out["error"]?.jsonObject ?: error("error block missing")
assertEquals("Unsupported field 'foo'", err["message"]?.jsonPrimitive?.content) assertEquals("Unsupported field 'foo'", err["message"]?.jsonPrimitive?.content)
} }
@Test
fun sha256MatchesKnownVectors() {
assertEquals(
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
sha256Hex(""),
)
assertEquals(
"ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad",
sha256Hex("abc"),
)
assertEquals(
"248d6a61d20638b8e5c026930c3e6039a33ce45964ff2167f6ecedd419db06c1",
sha256Hex("abcdbcdecdefdefgefghfghighijhijkijkljklmklmnlmnomnopnopq"),
)
}
@Test
fun sessionPrefixHashesAreStableWhenHistoryGrows() {
val first = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"S"},{"role":"user","content":"hi"}]}""",
).jsonObject
val second = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"S"},{"role":"user","content":"hi"},{"role":"assistant","content":"yo"},{"role":"user","content":"again"}]}""",
).jsonObject
val h1 = sessionPrefixHashes(first)
val h2 = sessionPrefixHashes(second)
// Цепочка начинается с первого user — system-преамбула не хэшируется.
assertEquals(1, h1.size)
assertEquals(3, h2.size)
assertEquals(h1, h2.take(1))
assertEquals(h1.last(), h2[0])
}
@Test
fun sessionPrefixHashesIgnoreLeadingSystemSoSharedPromptDoesNotCollide() {
val a = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"SAME"},{"role":"user","content":"session A"}]}""",
).jsonObject
val b = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"SAME"},{"role":"user","content":"session B"}]}""",
).jsonObject
// разные первые user-сообщения → разные хэши, несмотря на общий system
assertFalse(sessionPrefixHashes(a) == sessionPrefixHashes(b))
}
@Test
fun sessionPrefixHashesFallBackToWholeBodyWithoutMessages() {
val body = Json.parseToJsonElement("""{"model":"m"}""").jsonObject
assertEquals(listOf(sha256Hex(body.toString())), sessionPrefixHashes(body))
}
@Test
fun sessionRegistryReusesIdByLongestCommonPrefix() {
val reg = SessionRegistry()
val first = reg.resolveLocked(listOf("H1"), 0)
assertEquals("H1", first)
// история выросла: [H1, H2, H3] — самый длинный известный префикс H1
val next = reg.resolveLocked(listOf("H1", "H2", "H3"), 1)
assertEquals("H1", next)
// и дальше — id не меняется
val deep = reg.resolveLocked(listOf("H1", "H2", "H3", "H4"), 2)
assertEquals("H1", deep)
}
@Test
fun sessionRegistryCreatesNewIdForDifferentHistory() {
val reg = SessionRegistry()
assertEquals("A1", reg.resolveLocked(listOf("A1"), 0))
assertEquals("B1", reg.resolveLocked(listOf("B1"), 0))
assertEquals("A1", reg.resolveLocked(listOf("A1", "A2"), 0))
}
@Test
fun sessionRegistryEvictsLruWhenOverCapacity() {
val reg = SessionRegistry(maxSessions = 2, ttlMillis = Long.MAX_VALUE)
reg.resolveLocked(listOf("A1"), 0)
reg.resolveLocked(listOf("B1"), 1)
reg.resolveLocked(listOf("C1"), 2) // A вытеснена
// a1 больше неизвестен → новая сессия с id = хэш всей истории A2
assertEquals("A2", reg.resolveLocked(listOf("A1", "A2"), 3))
}
@Test
fun sessionRegistryEvictsExpiredByTtl() {
val reg = SessionRegistry(maxSessions = 100, ttlMillis = 1000)
reg.resolveLocked(listOf("A1"), 0)
reg.resolveLocked(listOf("B1"), 2000) // A протухла
assertEquals("A2", reg.resolveLocked(listOf("A1", "A2"), 2001))
}
@Test
fun parseConfigReadsSessionHeader() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
session_header: x-opencode-session
- id: p2
url: "https://y.ru/api/v1"
models:
- name: m1
upstreams: []
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals("x-opencode-session", cfg.providers[0].session_header)
assertEquals(null, cfg.providers[1].session_header)
}
@Test
fun parseConfigReadsThinkTagsOnProviderAndUpstream() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
think_tags: split
- id: p2
url: "https://y.ru/api/v1"
upstreams:
- id: u1
provider: p1
model: real-1
think_tags: strip
- id: u2
provider: p1
model: real-2
models:
- name: m1
upstreams: [u1, u2]
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals("split", cfg.providers[0].think_tags)
assertEquals(null, cfg.providers[1].think_tags)
assertEquals("strip", cfg.upstreams[0].think_tags)
assertEquals(null, cfg.upstreams[1].think_tags)
}
@Test
fun effectiveThinkTagsPrefersUpstreamThenProviderWithTolerantParse() {
val upSplit = UpstreamConf("u1", "p", "m", think_tags = "split")
val upNull = UpstreamConf("u2", "p", "m", think_tags = null)
val prov = { tt: String? -> ProviderConf("p", "https://x", think_tags = tt) }
// значение у апстрима — берётся оно, провайдер игнорируется
assertEquals("split", effectiveThinkTags(upSplit, null))
assertEquals("split", effectiveThinkTags(upSplit, prov("strip")))
// апстрим null — берётся провайдерский
assertEquals("split", effectiveThinkTags(upNull, prov("split")))
assertEquals("strip", effectiveThinkTags(upNull, prov("strip")))
assertEquals("split", effectiveThinkTags(upNull, prov("true")))
assertEquals("off", effectiveThinkTags(upNull, prov("false")))
assertEquals("off", effectiveThinkTags(upNull, prov("yes")))
// нигде нет — "off"
assertEquals("off", effectiveThinkTags(upNull, prov(null)))
assertEquals("off", effectiveThinkTags(upNull, null))
}
} }
@@ -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,75 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
class ThinkTagSplitterTest {
@Test
fun offIsPassthroughEvenWithTags() {
val s = ThinkTagSplitter("off")
assertEquals("<think>x</think>" to "", s.feed("<think>x</think>"))
assertEquals("привет\nещё" to "", s.feed("привет\nещё"))
assertEquals("" to "", s.finish())
}
@Test
fun splitWholeBlockMovesInnerTextToReasoning() {
val s = ThinkTagSplitter("split")
assertEquals("a" to "r", s.feed("<think>r</think>a"))
assertEquals("" to "", s.finish())
}
@Test
fun splitTagCutAcrossThreeChunks() {
val s = ThinkTagSplitter("split")
val first = s.feed("привет <t")
assertEquals("привет ", first.first)
assertEquals("", first.second)
assertEquals("" to "", s.feed("hi"))
val third = s.feed("nk>разум")
assertEquals("", third.first)
assertEquals("разум", third.second)
assertEquals("" to "", s.finish())
}
@Test
fun splitUnclosedOpenTagLeavesRestInReasoning() {
val s = ThinkTagSplitter("split")
assertEquals("" to "мысли без конца", s.feed("<think>мысли без конца"))
assertEquals("" to "", s.finish())
// хвост-префикс незакрытого закрывающего тега тоже уезжает в reasoning
val s2 = ThinkTagSplitter("split")
assertEquals("" to "abc", s2.feed("<think>abc</thi"))
assertEquals("" to "</thi", s2.finish())
}
@Test
fun splitHandlesSeveralBlocksInOneResponse() {
val s = ThinkTagSplitter("split")
val first = s.feed("<think>A</think>B<think>C</think>")
assertEquals("B", first.first)
assertEquals("AC", first.second)
assertEquals("text" to "", s.feed("text"))
assertEquals("" to "", s.finish())
}
@Test
fun stripDropsReasoningAndTags() {
val s = ThinkTagSplitter("strip")
assertEquals("y" to "", s.feed("<think>x</think>y"))
assertEquals("" to "", s.feed("<think>z"))
assertEquals("" to "", s.finish())
}
@Test
fun heldTailThatTurnedOutNotToBeATagFlowsBackAsContent() {
val s = ThinkTagSplitter("split")
// "<thi" — полный префикс открывающего тега: хвост держим
assertEquals("" to "", s.feed("<thi"))
// "x" не продолжает тег: хвост отдаём как обычный текст
assertEquals("<thix" to "", s.feed("x"))
assertEquals("" to "", s.finish())
}
}
@@ -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"))
}
}