Files
llm-proxy/docs/sse-done-contract.md
subochev 558b2f925b
Build LLM Proxy / Build and push (release) Successful in 38s
feat: выравнивание контракта конца SSE — дописывать data: [DONE] при штатном закрытии
Проблема (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

17 KiB
Raw Permalink Blame History

ТЗ: выравнивание контракта конца 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-хелперов)

/** Терминальный маркер 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(...)):

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)
}

Стало:

if (thinkMode == "off") {
    // Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
    ch.streamRawWithDoneContract(this)
} else {
    streamSseWithThinkTags(ch, thinkMode)
}

(если this в этом контексте — не ByteWriteChannel, использовать тот ресивер, который уже использовался существующими вызовами writeFully/flush в этой лямбде; компилятор — источник истины).

3.3. Правка streamSseWithThinkTags

Добавить два флага и запись фактов:

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. Команды проверки

./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 не трогаем: фикс на нашем хопе лечит всех клиентов сразу.