feat: флажок think_tags — вырезание think-тегов из ответа (split/strip)
Build LLM Proxy / Build and push (release) Successful in 24s
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.
This commit is contained in:
@@ -3,7 +3,9 @@ package pw.binom.llmproxy
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.call.body
|
||||
import io.ktor.client.request.headers
|
||||
import io.ktor.utils.io.LineEnding
|
||||
import io.ktor.utils.io.readAvailable
|
||||
import io.ktor.utils.io.readLine
|
||||
import io.ktor.http.Headers
|
||||
import io.ktor.client.request.preparePost
|
||||
import io.ktor.client.request.setBody
|
||||
@@ -236,16 +238,22 @@ private suspend fun handleChat(
|
||||
if (clientWantsStream) {
|
||||
val ct = resp.headers["Content-Type"] ?: "text/event-stream"
|
||||
val status = HttpStatusCode.fromValue(upstreamStatus)
|
||||
val thinkMode = effectiveThinkTags(up, provider)
|
||||
call.respondBytesWriter(ContentType.parse(ct), status) {
|
||||
val ch = resp.body<ByteReadChannel>()
|
||||
val buf = ByteArray(8192)
|
||||
while (true) {
|
||||
val n = ch.readAvailable(buf)
|
||||
if (n == -1) break
|
||||
if (n > 0) {
|
||||
writeFully(buf, 0, n)
|
||||
flush()
|
||||
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)
|
||||
}
|
||||
}
|
||||
log.info {
|
||||
@@ -255,7 +263,12 @@ private suspend fun handleChat(
|
||||
} else {
|
||||
val ct = resp.headers["Content-Type"] ?: "application/json"
|
||||
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 {
|
||||
if (ct.contains("text/event-stream") &&
|
||||
Json.parseToJsonElement(out).jsonObject["error"] != null
|
||||
@@ -416,6 +429,20 @@ internal fun tryClaim(up: UpstreamConf, active: Map<String, UpstreamCounter>): B
|
||||
internal fun effectiveConcurrencyLimit(up: UpstreamConf, provider: ProviderConf?): Int =
|
||||
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 по завершении проксирования). */
|
||||
internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) {
|
||||
active.getValue(up.id).release()
|
||||
@@ -580,6 +607,141 @@ internal fun rebuildFromChunks(sse: String): String {
|
||||
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-стрима с рассечением think-тегов. Строки, не начинающиеся
|
||||
* с `data:`, и `data: [DONE]` уходят клиенту без изменений (с `\n`). Прочие
|
||||
* `data:`-строки парсятся и прогоняются через [transformThinkChunk]; результат
|
||||
* записывается как `data: <json>\n\n` (событие-граница SSE), а при ошибке парса —
|
||||
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
|
||||
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
|
||||
*/
|
||||
private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
|
||||
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
|
||||
val addReasoning = thinkMode == "split"
|
||||
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")
|
||||
} else {
|
||||
val obj = runCatching { json.parseToJsonElement(payload).jsonObject }.getOrNull()
|
||||
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()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Пересобрать SSE-чанк: по каждому choice (ключ `index`, дефолт 0) взять
|
||||
* `delta.content` (только если это JSON-строка; массив частей не трогаем) и
|
||||
* прогнать через [ThinkTagSplitter] для этого index. Остаток возвращается в
|
||||
* `delta.content` (поле убирается, если пустое); вырезанное — в
|
||||
* `delta.reasoning_content` (только режим split, при strip не добавляем).
|
||||
* Чанк без `choices` или без строкового `delta.content` не меняется —
|
||||
* возвращается null (отдать исходную строку как есть).
|
||||
*/
|
||||
private 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
|
||||
data class ProviderConf(
|
||||
val id: String,
|
||||
@@ -588,6 +750,7 @@ data class ProviderConf(
|
||||
val max_concurrency: Int? = null,
|
||||
val patch: JsonObject? = null,
|
||||
val session_header: String? = null,
|
||||
val think_tags: String? = null,
|
||||
)
|
||||
|
||||
data class UpstreamConf(
|
||||
@@ -596,6 +759,7 @@ data class UpstreamConf(
|
||||
val model: String,
|
||||
val max_concurrency: Int? = null,
|
||||
val patch: JsonObject? = null,
|
||||
val think_tags: String? = null,
|
||||
)
|
||||
|
||||
data class ModelConf(
|
||||
@@ -650,6 +814,7 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
|
||||
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
|
||||
session_header = m.strOrNull("session_header"),
|
||||
think_tags = m.strOrNull("think_tags"),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -661,6 +826,7 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
model = m.str("model"),
|
||||
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
|
||||
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
|
||||
think_tags = m.strOrNull("think_tags"),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user