diff --git a/CONFIG.md b/CONFIG.md index 7417c4d..5656d3d 100644 --- a/CONFIG.md +++ b/CONFIG.md @@ -165,6 +165,59 @@ id не меняется; разные диалоги получают разн > истории рвёт общий префикс → сессия распадётся на новую. Диалоги с > одинаковым первым `user`-сообщением неразличимы (склеятся). +### Обработка think-тегов (`think_tags`) + +Некоторые провайдеры (например, **minimax**) отдают рассуждения модели не в +отдельном поле `reasoning_content`, а прямо в `content`, обернув их тегами +`…`. Флажок `think_tags` (опционально) разрешает прокси разрезать +такой ответ и разложить его по полям. + +Поле доступно на двух уровнях: + +- `providers[].think_tags` — правило по умолчанию для всех апстримов провайдера; +- `upstreams[].think_tags` — необязательное переопределение на конкретной + апстрим-модели. + +Значения (строки): + +| Значение | Поведение | +|---|---| +| `off` | дефолт: ответ не меняется, теги остаются в `content` | +| `split` | блоки `…` вырезаются из `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 — прокси держит хвост, который может оказаться +началом тега, и не отдаёт его клиенту до разрешения, поэтому огрызок тега не +утечёт. Незакрытый `` в конце потока трактуется как «всё после него — +рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ +не трогаем. Без флажка (`off`) ответ идёт байт-в-байт как раньше. + ### Пример сборки тела (многослойный `patch`) Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим diff --git a/README.md b/README.md index 03eb4f1..0e553ec 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,11 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` ( умолчанию для его апстримов), и на апстриме (перекрывает провайдерский); без обоих — безлимит. +Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или +апстрима (`off` по умолчанию, `split` — рассуждения из `…` уходят +в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в +non-stream. + | Переменная | Default | Описание | |---|---|---| | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt index 2bcccde..a988bf0 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt @@ -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() - 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() - 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): 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) { 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: \n\n` (событие-граница SSE), а при ошибке парса — + * исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в + * буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком. + */ +private suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) { + val splitters = mutableMapOf() + 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() + 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, + 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"), ) } diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/ThinkTagSplitter.kt b/src/commonMain/kotlin/pw/binom/llmproxy/ThinkTagSplitter.kt new file mode 100644 index 0000000..f708036 --- /dev/null +++ b/src/commonMain/kotlin/pw/binom/llmproxy/ThinkTagSplitter.kt @@ -0,0 +1,96 @@ +package pw.binom.llmproxy + +/** + * Автомат рассечения кусков стрима по think-тегам (строго lowercase). Чистая логика: без Ktor, без IO, без побочных + * эффектов — по одному экземпляру на choice/index. + * + * - `off` — passthrough: всё, включая теги, уходит в content; + * - `split` — текст внутри think-блока → reasoning, + * остальное → content; несколько блоков в одном ответе обрабатываются все; + * - `strip` — как `split`, но рассуждения выбрасываются (reasoning всегда ""). + * + * Тег может прийти разрезанным между кусками: хвост, который является + * префиксом ожидаемого тега: вне блока — , внутри блока — . + * Такой хвост не отдаём, держим до следующего `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 { + 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 { + 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 = "" + const val CLOSE = "" + } +} diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt index f020559..faf781d 100644 --- a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt +++ b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt @@ -494,4 +494,56 @@ class ConfigLogicTest { 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("true"))) + + // апстрим 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)) + } } diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/ThinkTagSplitterTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/ThinkTagSplitterTest.kt new file mode 100644 index 0000000..6b4f06b --- /dev/null +++ b/src/commonTest/kotlin/pw/binom/llmproxy/ThinkTagSplitterTest.kt @@ -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("x" to "", s.feed("x")) + assertEquals("привет\nещё" to "", s.feed("привет\nещё")) + assertEquals("" to "", s.finish()) + } + + @Test + fun splitWholeBlockMovesInnerTextToReasoning() { + val s = ThinkTagSplitter("split") + assertEquals("a" to "r", s.feed("ra")) + assertEquals("" to "", s.finish()) + } + + @Test + fun splitTagCutAcrossThreeChunks() { + val s = ThinkTagSplitter("split") + val first = s.feed("привет разум") + assertEquals("", third.first) + assertEquals("разум", third.second) + assertEquals("" to "", s.finish()) + } + + @Test + fun splitUnclosedOpenTagLeavesRestInReasoning() { + val s = ThinkTagSplitter("split") + assertEquals("" to "мысли без конца", s.feed("мысли без конца")) + assertEquals("" to "", s.finish()) + + // хвост-префикс незакрытого закрывающего тега тоже уезжает в reasoning + val s2 = ThinkTagSplitter("split") + assertEquals("" to "abc", s2.feed("abcABC") + 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("xy")) + assertEquals("" to "", s.feed("z")) + assertEquals("" to "", s.finish()) + } + + @Test + fun heldTailThatTurnedOutNotToBeATagFlowsBackAsContent() { + val s = ThinkTagSplitter("split") + // "