From 558b2f925bc21c7457ec2c5af1b0551901bcf298 Mon Sep 17 00:00:00 2001 From: Hermes Agent Date: Fri, 11 Sep 2026 18:26:52 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20=D0=B2=D1=8B=D1=80=D0=B0=D0=B2=D0=BD?= =?UTF-8?q?=D0=B8=D0=B2=D0=B0=D0=BD=D0=B8=D0=B5=20=D0=BA=D0=BE=D0=BD=D1=82?= =?UTF-8?q?=D1=80=D0=B0=D0=BA=D1=82=D0=B0=20=D0=BA=D0=BE=D0=BD=D1=86=D0=B0?= =?UTF-8?q?=20SSE=20=E2=80=94=20=D0=B4=D0=BE=D0=BF=D0=B8=D1=81=D1=8B=D0=B2?= =?UTF-8?q?=D0=B0=D1=82=D1=8C=20data:=20[DONE]=20=D0=BF=D1=80=D0=B8=20?= =?UTF-8?q?=D1=88=D1=82=D0=B0=D1=82=D0=BD=D0=BE=D0=BC=20=D0=B7=D0=B0=D0=BA?= =?UTF-8?q?=D1=80=D1=8B=D1=82=D0=B8=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Проблема (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. --- TESTING.md | 6 +- docs/sse-done-contract.md | 280 ++++++++++++++++++ .../kotlin/pw/binom/llmproxy/Main.kt | 111 ++++++- .../binom/llmproxy/StreamDoneContractTest.kt | 176 +++++++++++ 4 files changed, 561 insertions(+), 12 deletions(-) create mode 100644 docs/sse-done-contract.md create mode 100644 src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt diff --git a/TESTING.md b/TESTING.md index c669e24..45505da 100644 --- a/TESTING.md +++ b/TESTING.md @@ -17,6 +17,7 @@ | `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`, отсутствие дублирования. | ## Проверка качества тестов (мутационная приёмка) @@ -36,7 +37,10 @@ - `transformThinkChunk` возвращает null → падают тесты чанков и стрима; - блок финиш-чанка в `streamSseWithThinkTags` не выполняется → падает тест про удержанный хвост; - `holdableSuffix` всегда 0 (хвост тега не удерживается) → падают тесты автомата, чанков и стрима; -- `strip` начинает отдавать рассуждения → падает тест автомата. +- `strip` начинает отдавать рассуждения → падает тест автомата; +- `missingDoneMarker` всегда `false` → падают тесты сырого пути (`rawStreamIsByteExactAndAppendsDone`) и тест разрезанного маркера; +- `if (sawFinishReason && !sawDone)` → `if (!sawDone)` в `streamSseWithThinkTags` → падает `thinkStreamTruncatedNeedsNoMarker`; +- `if (sawFinishReason && !sawDone)` → `if (false)` в `streamSseWithThinkTags` → падает `thinkStreamAppendsDoneWhenUpstreamClosedWithoutMarker`. Правило: боевой код нельзя подгонять под тест; если тест не проходит, неверен тест. diff --git a/docs/sse-done-contract.md b/docs/sse-done-contract.md new file mode 100644 index 0000000..f17ec4d --- /dev/null +++ b/docs/sse-done-contract.md @@ -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() + 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** не трогаем: фикс на нашем хопе лечит всех клиентов сразу. diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt index 4daa380..0aff6dc 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt @@ -242,16 +242,8 @@ private suspend fun handleChat( call.respondBytesWriter(ContentType.parse(ct), status) { val ch = resp.body() if (thinkMode == "off") { - // Флажок не выставлен — сырой байтовый passthrough как раньше. - val buf = ByteArray(8192) - while (true) { - val n = ch.readAvailable(buf) - if (n == -1) break - if (n > 0) { - writeFully(buf, 0, n) - flush() - } - } + // Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока. + streamRawWithDoneContract(ch) } else { streamSseWithThinkTags(ch, thinkMode) } @@ -644,6 +636,93 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String? 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`). Прочие @@ -655,15 +734,19 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String? internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) { val splitters = mutableMapOf() val addReasoning = thinkMode == "split" + var sawDone = false + var sawFinishReason = false while (true) { val line = source.readLine(LineEnding.Lenient) ?: break when { line.startsWith("data:") -> { val payload = line.removePrefix("data:").trim() if (payload == "[DONE]") { - emitUtf8("data: [DONE]\n\n") + sawDone = true + emitUtf8(SSE_DONE_MARKER) } else { val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull() + if (obj != null && hasFinishReason(obj)) sawFinishReason = true val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) } if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n") } @@ -699,6 +782,12 @@ internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadCha flush() } } + // После сброса хвостов дописываем маркер, только если апстрим его не прислал, + // а `finish_reason` в потоке был (при обрыве без него маркер НЕ дописываем). + if (sawFinishReason && !sawDone) { + emitUtf8(SSE_DONE_MARKER) + flush() + } } /** diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt new file mode 100644 index 0000000..26f6006 --- /dev/null +++ b/src/commonTest/kotlin/pw/binom/llmproxy/StreamDoneContractTest.kt @@ -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), "маркер не должен дублироваться") + } +}