feat: выравнивание контракта конца SSE — дописывать data: [DONE] при штатном закрытии
Build LLM Proxy / Build and push (release) Successful in 38s
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.
This commit is contained in:
@@ -242,16 +242,8 @@ private suspend fun handleChat(
|
||||
call.respondBytesWriter(ContentType.parse(ct), status) {
|
||||
val ch = resp.body<ByteReadChannel>()
|
||||
if (thinkMode == "off") {
|
||||
// Флажок не выставлен — сырой байтовый passthrough как раньше.
|
||||
val buf = ByteArray(8192)
|
||||
while (true) {
|
||||
val n = ch.readAvailable(buf)
|
||||
if (n == -1) break
|
||||
if (n > 0) {
|
||||
writeFully(buf, 0, n)
|
||||
flush()
|
||||
}
|
||||
}
|
||||
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
|
||||
streamRawWithDoneContract(ch)
|
||||
} else {
|
||||
streamSseWithThinkTags(ch, thinkMode)
|
||||
}
|
||||
@@ -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<Int, ThinkTagSplitter>()
|
||||
val addReasoning = thinkMode == "split"
|
||||
var sawDone = false
|
||||
var sawFinishReason = false
|
||||
while (true) {
|
||||
val line = source.readLine(LineEnding.Lenient) ?: break
|
||||
when {
|
||||
line.startsWith("data:") -> {
|
||||
val payload = line.removePrefix("data:").trim()
|
||||
if (payload == "[DONE]") {
|
||||
emitUtf8("data: [DONE]\n\n")
|
||||
sawDone = true
|
||||
emitUtf8(SSE_DONE_MARKER)
|
||||
} else {
|
||||
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
|
||||
if (obj != null && hasFinishReason(obj)) sawFinishReason = true
|
||||
val out = obj?.let { transformThinkChunk(it, splitters, thinkMode, addReasoning) }
|
||||
if (out == null) emitUtf8("$line\n\n") else emitUtf8("data: $out\n\n")
|
||||
}
|
||||
@@ -699,6 +782,12 @@ internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadCha
|
||||
flush()
|
||||
}
|
||||
}
|
||||
// После сброса хвостов дописываем маркер, только если апстрим его не прислал,
|
||||
// а `finish_reason` в потоке был (при обрыве без него маркер НЕ дописываем).
|
||||
if (sawFinishReason && !sawDone) {
|
||||
emitUtf8(SSE_DONE_MARKER)
|
||||
flush()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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), "маркер не должен дублироваться")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user