3 Commits
7 .. 8

Author SHA1 Message Date
subochev c7e7c7346d feat: флажок think_tags — вырезание think-тегов из ответа (split/strip)
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.
2026-09-11 16:24:54 +03:00
subochev dba52d22fd feat: вычисление сессии по истории (LCP) для providers[].session_header
- SHA-256 на commonMain + SessionRegistry (LCP по префикс-хэшам, LRU 1000/6ч, Mutex)
- providers[].session_header: прокси сам ставит/перезаписывает заголовок сессии
- цепочка хэшей стартует с первого user-сообщения (system не склеивает сессии)
- CONFIG.md + тесты (SHA-256, LCP, LRU/TTL, парсинг, стабильность префиксов)
2026-09-11 05:27:15 +03:00
subochev efdce75ee9 feat: логирование заголовков запроса и заголовков в апстрим; маскировка секретов; Content-Type прокси-овнер 2026-09-11 04:44:05 +03:00
7 changed files with 891 additions and 11 deletions
+88
View File
@@ -73,6 +73,7 @@ providers:
url: "https://routerai.ru/api/v1" url: "https://routerai.ru/api/v1"
key: "sk-..." # Bearer-ключ; можно подставлять из env key: "sk-..." # Bearer-ключ; можно подставлять из env
max_concurrency: 4 # опционально; лимит по умолчанию для апстримов max_concurrency: 4 # опционально; лимит по умолчанию для апстримов
session_header: x-opencode-session # опционально; прокси считает сессию из истории
patch: # уровень провайдера: ко всем его запросам patch: # уровень провайдера: ко всем его запросам
provider: provider:
allow_fallbacks: false allow_fallbacks: false
@@ -131,6 +132,92 @@ models:
`patch` опционален на **любом** уровне (`providers` / `upstreams` / `models`): `patch` опционален на **любом** уровне (`providers` / `upstreams` / `models`):
если ни одного нет — запрос проксируется как есть (исходное тело клиента). если ни одного нет — запрос проксируется как есть (исходное тело клиента).
### Сессия по истории (`session_header`)
`providers[].session_header` (опционально) — имя HTTP-заголовка, который прокси
**вычисляет сам** из истории сообщений и ставит в запрос к этому провайдеру.
Нужно для API, требующих стабильный идентификатор сессии (например,
`x-opencode-session`), когда клиент его не шлёт или шлёт не то.
```yaml
providers:
- id: some-provider
url: "https://.../v1"
session_header: x-opencode-session
```
Как считается id:
1. Берётся финальное тело запроса (после всех `patch`), из него — `messages`.
2. Цепочка **инкрементальных SHA-256 префикс-хэшей** начинается с первого
`user`-сообщения (ведущий `system`-промпт игнорируется: он обычно одинаков
у всех сессий клиента и как признак сессии бесполезен).
3. В реестре сессий ищется **наибольший общий префикс** (LCP) с уже виденной
историей. Нашли — используется id той сессии; не нашли — создаётся новая
(`id` = хэш всей истории на первом ходу).
4. Заголовок ставится **всегда** (клиентское значение перезаписывается).
Итог: пока история одной сессии растёт (дописываются assistant/user-сообщения),
id не меняется; разные диалоги получают разные id.
> **Ограничения.** Реестр живёт в памяти (LRU: 1000 сессий / 6 часов) — при
> рестарте прокси активные сессии получат новый id. Обрезка/суммаризация
> истории рвёт общий префикс → сессия распадётся на новую. Диалоги с
> одинаковым первым `user`-сообщением неразличимы (склеятся).
### Обработка think-тегов (`think_tags`)
Некоторые провайдеры (например, **minimax**) отдают рассуждения модели не в
отдельном поле `reasoning_content`, а прямо в `content`, обернув их тегами
`<think>…</think>`. Флажок `think_tags` (опционально) разрешает прокси разрезать
такой ответ и разложить его по полям.
Поле доступно на двух уровнях:
- `providers[].think_tags` — правило по умолчанию для всех апстримов провайдера;
- `upstreams[].think_tags` — необязательное переопределение на конкретной
апстрим-модели.
Значения (строки):
| Значение | Поведение |
|---|---|
| `off` | дефолт: ответ не меняется, теги остаются в `content` |
| `split` | блоки `<think>…</think>` вырезаются из `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 — прокси держит хвост, который может оказаться
началом тега, и не отдаёт его клиенту до разрешения, поэтому огрызок тега не
утечёт. Незакрытый `<think>` в конце потока трактуется как «всё после него —
рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ
не трогаем. Без флажка (`off`) ответ идёт байт-в-байт как раньше.
### Пример сборки тела (многослойный `patch`) ### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -211,6 +298,7 @@ data class ProviderConf(
val key: String = "", val key: String = "",
val max_concurrency: Int? = null, // лимит по умолчанию для апстримов провайдера val max_concurrency: Int? = null, // лимит по умолчанию для апстримов провайдера
val patch: JsonObject? = null, // ко всем запросам провайдера val patch: JsonObject? = null, // ко всем запросам провайдера
val session_header: String? = null, // заголовок-сессия, считается из истории
) )
@Serializable @Serializable
+5
View File
@@ -19,6 +19,11 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` (
умолчанию для его апстримов), и на апстриме (перекрывает провайдерский); без умолчанию для его апстримов), и на апстриме (перекрывает провайдерский); без
обоих — безлимит. обоих — безлимит.
Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или
апстрима (`off` по умолчанию, `split` — рассуждения из `<think>…</think>` уходят
в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в
non-stream.
| Переменная | Default | Описание | | Переменная | Default | Описание |
|---|---|---| |---|---|---|
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
+222 -4
View File
@@ -3,7 +3,9 @@ package pw.binom.llmproxy
import io.ktor.client.HttpClient import io.ktor.client.HttpClient
import io.ktor.client.call.body import io.ktor.client.call.body
import io.ktor.client.request.headers import io.ktor.client.request.headers
import io.ktor.utils.io.LineEnding
import io.ktor.utils.io.readAvailable import io.ktor.utils.io.readAvailable
import io.ktor.utils.io.readLine
import io.ktor.http.Headers import io.ktor.http.Headers
import io.ktor.client.request.preparePost import io.ktor.client.request.preparePost
import io.ktor.client.request.setBody import io.ktor.client.request.setBody
@@ -14,7 +16,9 @@ import io.ktor.server.application.Application
import io.ktor.server.application.ApplicationCall import io.ktor.server.application.ApplicationCall
import io.ktor.server.application.call import io.ktor.server.application.call
import io.ktor.server.application.install import io.ktor.server.application.install
import io.ktor.server.request.httpMethod
import io.ktor.server.request.receiveText import io.ktor.server.request.receiveText
import io.ktor.server.request.uri
import io.ktor.server.response.respondBytesWriter import io.ktor.server.response.respondBytesWriter
import io.ktor.server.response.respondText import io.ktor.server.response.respondText
import io.ktor.server.routing.get import io.ktor.server.routing.get
@@ -95,8 +99,9 @@ fun main() {
} }
val http = createHttpClient() val http = createHttpClient()
val sessions = SessionRegistry()
startServer(config.server.host, config.server.port) { startServer(config.server.host, config.server.port) {
proxyModule(config, providersById, upstreamsById, active, http) proxyModule(config, providersById, upstreamsById, active, sessions, http)
} }
} }
@@ -107,11 +112,12 @@ fun Application.proxyModule(
providersById: Map<String, ProviderConf>, providersById: Map<String, ProviderConf>,
upstreamsById: Map<String, UpstreamConf>, upstreamsById: Map<String, UpstreamConf>,
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
) { ) {
routing { routing {
post("/v1/chat/completions") { post("/v1/chat/completions") {
handleChat(call, config, providersById, upstreamsById, active, http) handleChat(call, config, providersById, upstreamsById, active, sessions, http)
} }
get("/v1/models") { get("/v1/models") {
handleModels(call, config) handleModels(call, config)
@@ -125,8 +131,10 @@ private suspend fun handleChat(
providersById: Map<String, ProviderConf>, providersById: Map<String, ProviderConf>,
upstreamsById: Map<String, UpstreamConf>, upstreamsById: Map<String, UpstreamConf>,
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
) { ) {
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val raw = call.receiveText() val raw = call.receiveText()
if (raw.isBlank()) { if (raw.isBlank()) {
call.respondText(errorJson("empty body"), ContentType.Application.Json, HttpStatusCode.BadRequest) call.respondText(errorJson("empty body"), ContentType.Application.Json, HttpStatusCode.BadRequest)
@@ -183,13 +191,27 @@ private suspend fun handleChat(
val url = provider.url.trimEnd('/') + "/chat/completions" val url = provider.url.trimEnd('/') + "/chat/completions"
val providerKey = resolveEnv(provider.key) val providerKey = resolveEnv(provider.key)
val sessionHeader = provider.session_header
// Клиентский одноимённый заголовок не пробрасываем: при метке
// значение x-opencode-session всегда вычисляем сами (LCP по истории).
val forwardedHeaders = headersToForward(call.request.headers)
.filterKeys { sessionHeader == null || !it.equals(sessionHeader, ignoreCase = true) }
val sessionId = sessionHeader?.let { sessions.resolve(sessionPrefixHashes(forwarded)) }
val outgoingHeaders = forwardedHeaders.toMutableMap().apply {
this["Content-Type"] = listOf("application/json")
if (providerKey.isNotEmpty()) this["Authorization"] = listOf("Bearer $providerKey")
if (sessionHeader != null && sessionId != null) this[sessionHeader] = listOf(sessionId)
}
log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} session=${sessionId ?: "-"} → $url headers: ${formatHeadersForLog(outgoingHeaders)}"
}
var failover = false var failover = false
var responded = false var responded = false
var upstreamStatus = 0 var upstreamStatus = 0
http.preparePost(url) { http.preparePost(url) {
headers { headers {
headersToForward(call.request.headers).forEach { (name, values) -> forwardedHeaders.forEach { (name, values) ->
appendAll(name, values) appendAll(name, values)
} }
// Авторизация — всегда наша (ключ провайдера из конфига); // Авторизация — всегда наша (ключ провайдера из конфига);
@@ -199,6 +221,9 @@ private suspend fun handleChat(
set("Authorization", "Bearer $providerKey") set("Authorization", "Bearer $providerKey")
} }
set("Content-Type", "application/json") set("Content-Type", "application/json")
if (sessionHeader != null && sessionId != null) {
set(sessionHeader, sessionId)
}
} }
setBody(forwarded.toString()) setBody(forwarded.toString())
}.execute { resp -> }.execute { resp ->
@@ -213,8 +238,11 @@ private suspend fun handleChat(
if (clientWantsStream) { if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream" val ct = resp.headers["Content-Type"] ?: "text/event-stream"
val status = HttpStatusCode.fromValue(upstreamStatus) val status = HttpStatusCode.fromValue(upstreamStatus)
val thinkMode = effectiveThinkTags(up, provider)
call.respondBytesWriter(ContentType.parse(ct), status) { call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>() val ch = resp.body<ByteReadChannel>()
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough как раньше.
val buf = ByteArray(8192) val buf = ByteArray(8192)
while (true) { while (true) {
val n = ch.readAvailable(buf) val n = ch.readAvailable(buf)
@@ -224,6 +252,9 @@ private suspend fun handleChat(
flush() flush()
} }
} }
} else {
streamSseWithThinkTags(ch, thinkMode)
}
} }
log.info { log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " + "[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
@@ -232,7 +263,12 @@ private suspend fun handleChat(
} else { } else {
val ct = resp.headers["Content-Type"] ?: "application/json" val ct = resp.headers["Content-Type"] ?: "application/json"
val full = resp.body<String>() 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 { val outCt = runCatching {
if (ct.contains("text/event-stream") && if (ct.contains("text/event-stream") &&
Json.parseToJsonElement(out).jsonObject["error"] != null Json.parseToJsonElement(out).jsonObject["error"] != null
@@ -287,6 +323,9 @@ private val SKIP_HEADER_NAMES = setOf(
// Authorization управляется прокси явно (ключ провайдера), клиентский // Authorization управляется прокси явно (ключ провайдера), клиентский
// не пересылается // не пересылается
"authorization", "authorization",
// Content-Type всегда наш (application/json: тело мержится как JSON),
// клиентский не пересылаем, чтобы не ушло двух заголовков
"content-type",
) )
/** /**
@@ -302,6 +341,29 @@ internal fun headersToForward(request: Headers): Map<String, List<String>> =
.filter { (name, _) -> name.lowercase() !in SKIP_HEADER_NAMES } .filter { (name, _) -> name.lowercase() !in SKIP_HEADER_NAMES }
.associate { (name, values) -> name to values } .associate { (name, values) -> name to values }
/** Заголовки, значения которых маскируются в логах (секреты клиента). */
private val SENSITIVE_HEADER_NAMES = setOf(
"authorization",
"proxy-authorization",
"x-api-key",
"api-key",
"cookie",
"set-cookie",
)
/**
* Заголовки в виде строки для лога (`name=v1|v2, ...`). Значения чувствительных
* имён ([SENSITIVE_HEADER_NAMES]) маскируются `***`, чтобы не светить секреты.
*/
internal fun formatHeadersForLog(headers: Map<String, List<String>>): String =
headers.entries.joinToString(", ") { (name, values) ->
val shown = if (name.lowercase() in SENSITIVE_HEADER_NAMES) values.map { "***" } else values
"$name=${shown.joinToString("|")}"
}
internal fun formatHeadersForLog(headers: Headers): String =
formatHeadersForLog(headers.entries().associate { (name, values) -> name to values })
/** /**
* Сборка тела запроса: подмена `model` на реальное имя апстрима + глубокий * Сборка тела запроса: подмена `model` на реальное имя апстрима + глубокий
* послойный мерж `patch` в порядке provider → upstream → model. * послойный мерж `patch` в порядке provider → upstream → model.
@@ -367,6 +429,20 @@ internal fun tryClaim(up: UpstreamConf, active: Map<String, UpstreamCounter>): B
internal fun effectiveConcurrencyLimit(up: UpstreamConf, provider: ProviderConf?): Int = internal fun effectiveConcurrencyLimit(up: UpstreamConf, provider: ProviderConf?): Int =
up.max_concurrency ?: provider?.max_concurrency ?: Int.MAX_VALUE 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 по завершении проксирования). */ /** Освободить слот апстрима (в finally по завершении проксирования). */
internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) { internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) {
active.getValue(up.id).release() active.getValue(up.id).release()
@@ -422,6 +498,7 @@ internal fun pickFreeUpstream(
pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) } pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) }
private suspend fun handleModels(call: ApplicationCall, config: Config) { private suspend fun handleModels(call: ApplicationCall, config: Config) {
log.info { "[llm-proxy] models ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val created = TimeSource.Monotonic.markNow().elapsedNow().inWholeSeconds val created = TimeSource.Monotonic.markNow().elapsedNow().inWholeSeconds
val data = config.models.map { m -> val data = config.models.map { m ->
JsonObject( JsonObject(
@@ -530,6 +607,141 @@ internal fun rebuildFromChunks(sse: String): String {
return JsonObject(root).toString() 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 @Serializable
data class ProviderConf( data class ProviderConf(
val id: String, val id: String,
@@ -537,6 +749,8 @@ data class ProviderConf(
val key: String = "", val key: String = "",
val max_concurrency: Int? = null, val max_concurrency: Int? = null,
val patch: JsonObject? = null, val patch: JsonObject? = null,
val session_header: String? = null,
val think_tags: String? = null,
) )
data class UpstreamConf( data class UpstreamConf(
@@ -545,6 +759,7 @@ data class UpstreamConf(
val model: String, val model: String,
val max_concurrency: Int? = null, val max_concurrency: Int? = null,
val patch: JsonObject? = null, val patch: JsonObject? = null,
val think_tags: String? = null,
) )
data class ModelConf( data class ModelConf(
@@ -598,6 +813,8 @@ internal fun parseConfig(root: YamlElement): Config {
key = m.strOrNull("key") ?: "", key = m.strOrNull("key") ?: "",
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
session_header = m.strOrNull("session_header"),
think_tags = m.strOrNull("think_tags"),
) )
} }
@@ -609,6 +826,7 @@ internal fun parseConfig(root: YamlElement): Config {
model = m.str("model"), model = m.str("model"),
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
think_tags = m.strOrNull("think_tags"),
) )
} }
@@ -0,0 +1,206 @@
package pw.binom.llmproxy
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.datetime.Clock
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonElement
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive
/**
* SHA-256, чистая реализация на commonMain (без внешних зависимостей —
* KMP-зеркало их не отдаёт). Нужен для стабильного хэша префиксов истории.
*/
internal object Sha256 {
private val K = intArrayOf(
0x428a2f98, 0x71374491, 0xb5c0fbcf.toInt(), 0xe9b5dba5.toInt(), 0x3956c25b, 0x59f111f1, 0x923f82a4.toInt(), 0xab1c5ed5.toInt(),
0xd807aa98.toInt(), 0x12835b01, 0x243185be, 0x550c7dc3, 0x72be5d74, 0x80deb1fe.toInt(), 0x9bdc06a7.toInt(), 0xc19bf174.toInt(),
0xe49b69c1.toInt(), 0xefbe4786.toInt(), 0x0fc19dc6, 0x240ca1cc, 0x2de92c6f, 0x4a7484aa, 0x5cb0a9dc, 0x76f988da,
0x983e5152.toInt(), 0xa831c66d.toInt(), 0xb00327c8.toInt(), 0xbf597fc7.toInt(), 0xc6e00bf3.toInt(), 0xd5a79147.toInt(), 0x06ca6351, 0x14292967,
0x27b70a85, 0x2e1b2138, 0x4d2c6dfc, 0x53380d13, 0x650a7354, 0x766a0abb, 0x81c2c92e.toInt(), 0x92722c85.toInt(),
0xa2bfe8a1.toInt(), 0xa81a664b.toInt(), 0xc24b8b70.toInt(), 0xc76c51a3.toInt(), 0xd192e819.toInt(), 0xd6990624.toInt(), 0xf40e3585.toInt(), 0x106aa070,
0x19a4c116, 0x1e376c08, 0x2748774c, 0x34b0bcb5, 0x391c0cb3, 0x4ed8aa4a, 0x5b9cca4f, 0x682e6ff3,
0x748f82ee, 0x78a5636f, 0x84c87814.toInt(), 0x8cc70208.toInt(), 0x90befffa.toInt(), 0xa4506ceb.toInt(), 0xbef9a3f7.toInt(), 0xc67178f2.toInt(),
)
private fun rotr(x: Int, n: Int): Int = (x ushr n) or (x shl (32 - n))
private fun pad(input: ByteArray): ByteArray {
val total = ((input.size + 9 + 63) / 64) * 64
val out = ByteArray(total)
input.copyInto(out)
out[input.size] = 0x80.toByte()
val bits = input.size.toLong() * 8
for (k in 0 until 8) {
out[total - 1 - k] = (bits ushr (8 * k)).toByte()
}
return out
}
fun hash(input: ByteArray): ByteArray {
val msg = pad(input)
val h = intArrayOf(
0x6a09e667, 0xbb67ae85.toInt(), 0x3c6ef372, 0xa54ff53a.toInt(),
0x510e527f, 0x9b05688c.toInt(), 0x1f83d9ab, 0x5be0cd19,
)
val w = IntArray(64)
var i = 0
while (i < msg.size) {
for (j in 0 until 16) {
w[j] = ((msg[i + j * 4].toInt() and 0xff) shl 24) or
((msg[i + j * 4 + 1].toInt() and 0xff) shl 16) or
((msg[i + j * 4 + 2].toInt() and 0xff) shl 8) or
(msg[i + j * 4 + 3].toInt() and 0xff)
}
for (j in 16 until 64) {
val s0 = rotr(w[j - 15], 7) xor rotr(w[j - 15], 18) xor (w[j - 15] ushr 3)
val s1 = rotr(w[j - 2], 17) xor rotr(w[j - 2], 19) xor (w[j - 2] ushr 10)
w[j] = w[j - 16] + s0 + w[j - 7] + s1
}
var a = h[0]; var b = h[1]; var c = h[2]; var d = h[3]
var e = h[4]; var f = h[5]; var g = h[6]; var hh = h[7]
for (j in 0 until 64) {
val s1 = rotr(e, 6) xor rotr(e, 11) xor rotr(e, 25)
val ch = (e and f) xor (e.inv() and g)
val t1 = hh + s1 + ch + K[j] + w[j]
val s0 = rotr(a, 2) xor rotr(a, 13) xor rotr(a, 22)
val maj = (a and b) xor (a and c) xor (b and c)
val t2 = s0 + maj
hh = g; g = f; f = e; e = d + t1; d = c; c = b; b = a; a = t1 + t2
}
h[0] += a; h[1] += b; h[2] += c; h[3] += d
h[4] += e; h[5] += f; h[6] += g; h[7] += hh
i += 64
}
val out = ByteArray(32)
for (j in 0 until 8) {
out[j * 4] = (h[j] ushr 24).toByte()
out[j * 4 + 1] = (h[j] ushr 16).toByte()
out[j * 4 + 2] = (h[j] ushr 8).toByte()
out[j * 4 + 3] = h[j].toByte()
}
return out
}
}
private const val HEX = "0123456789abcdef"
/** SHA-256 строки (UTF-8) в нижнем hex. */
internal fun sha256Hex(text: String): String {
val bytes = Sha256.hash(text.encodeToByteArray())
val sb = StringBuilder(bytes.size * 2)
for (b in bytes) {
val v = b.toInt() and 0xff
sb.append(HEX[v ushr 4]).append(HEX[v and 0xf])
}
return sb.toString()
}
/**
* Инкрементальные префикс-хэши истории: H_i = sha256(H_{i-1} + "\u0000" + msg_i).
*
* Цепочка начинается с первого `user`-сообщения: system-промпт обычно
* идентичен у всех сессий одного клиента и как признак сессии бесполезен
* (иначе LCP склеивает все сессии на общем префиксе `[system]`).
* Если `messages` нет/пусто — fallback на хэш всего тела.
*/
internal fun sessionPrefixHashes(body: JsonObject): List<String> {
val messages = body["messages"] as? JsonArray
if (messages == null || messages.isEmpty()) {
return listOf(sha256Hex(body.toString()))
}
fun roleAt(i: Int): String? =
((messages[i] as? JsonObject)?.get("role") as? JsonPrimitive)?.content
var start = messages.indices.firstOrNull { roleAt(it) == "user" }
?: messages.indices.firstOrNull { roleAt(it) != "system" }
?: 0
val res = ArrayList<String>(messages.size - start)
var prev = ""
for (i in start until messages.size) {
prev = sha256Hex(prev + "\u0000" + messages[i].toString())
res.add(prev)
}
return res
}
/**
* Реестр сессий: id сессии определяется наибольшим общим префиксом (LCP)
* присланной истории. Для префиксов, которые уже встречались, возвращается
* id исходной сессии; иначе создаётся новая (id = хэш всей истории).
*
* Таблица ограничена по размеру (LRU) и времени жизни (TTL). Доступ под
* Mutex: параллельные запросы одной сессии не должны гонять состояние.
*/
class SessionRegistry(
private val maxSessions: Int = 1000,
private val ttlMillis: Long = 6 * 60 * 60 * 1000L,
) {
private class Entry(val prefixes: MutableSet<String>) {
var lastAccess: Long = 0
}
private val mutex = Mutex()
private val prefixToSession = mutableMapOf<String, String>()
/** В порядке доступа: голова — самая давняя, хвост — свежая. */
private val sessions = LinkedHashMap<String, Entry>()
suspend fun resolve(prefixHashes: List<String>): String? =
mutex.withLock { resolveLocked(prefixHashes, Clock.System.now().toEpochMilliseconds()) }
/** Синхронная (без блокировки) версия — для тестов и вызовов под mutex. */
internal fun resolveLocked(prefixHashes: List<String>, now: Long): String? {
if (prefixHashes.isEmpty()) return null
evictExpired(now)
var found: String? = null
for (i in prefixHashes.indices.reversed()) {
val s = prefixToSession[prefixHashes[i]]
if (s != null) {
found = s
break
}
}
val id = found ?: prefixHashes.last()
val entry = sessions.remove(id) ?: Entry(mutableSetOf())
for (h in prefixHashes) {
val prev = prefixToSession.put(h, id)
if (prev != null && prev != id) {
sessions[prev]?.prefixes?.remove(h)
}
entry.prefixes.add(h)
}
entry.lastAccess = now
sessions[id] = entry
evictOverflow()
return id
}
private fun drop(sessionId: String, entry: Entry) {
entry.prefixes.forEach { h -> if (prefixToSession[h] == sessionId) prefixToSession.remove(h) }
}
private fun evictExpired(now: Long) {
val it = sessions.entries.iterator()
while (it.hasNext()) {
val e = it.next()
if (now - e.value.lastAccess <= ttlMillis) break
drop(e.key, e.value)
it.remove()
}
}
private fun evictOverflow() {
val it = sessions.entries.iterator()
while (sessions.size > maxSessions && it.hasNext()) {
val e = it.next()
drop(e.key, e.value)
it.remove()
}
}
}
@@ -0,0 +1,96 @@
package pw.binom.llmproxy
/**
* Автомат рассечения кусков стрима по think-тегам (строго lowercase). Чистая логика: без Ktor, без IO, без побочных
* эффектов — по одному экземпляру на choice/index.
*
* - `off` — passthrough: всё, включая теги, уходит в content;
* - `split` — текст внутри think-блока → reasoning,
* остальное → content; несколько блоков в одном ответе обрабатываются все;
* - `strip` — как `split`, но рассуждения выбрасываются (reasoning всегда "").
*
* Тег может прийти разрезанным между кусками: хвост, который является
* префиксом ожидаемого тега: вне блока — <think>, внутри блока — </think>.
* Такой хвост не отдаём, держим до следующего `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<String, String> {
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<String, String> {
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 = "<think>"
const val CLOSE = "</think>"
}
}
@@ -289,6 +289,7 @@ class ConfigLogicTest {
"Proxy-Connection" to listOf("keep-alive"), "Proxy-Connection" to listOf("keep-alive"),
"Upgrade" to listOf("h2c"), "Upgrade" to listOf("h2c"),
"Authorization" to listOf("Bearer client-secret"), "Authorization" to listOf("Bearer client-secret"),
"Content-Type" to listOf("application/x-www-form-urlencoded"),
"Accept" to listOf("*/*"), "Accept" to listOf("*/*"),
) )
val out = headersToForward(req) val out = headersToForward(req)
@@ -297,6 +298,7 @@ class ConfigLogicTest {
assertEquals(listOf("*/*"), out["Accept"]) assertEquals(listOf("*/*"), out["Accept"])
assertEquals(3, out.size) assertEquals(3, out.size)
assertEquals(null, out["Authorization"]) assertEquals(null, out["Authorization"])
assertEquals(null, out["Content-Type"])
} }
@Test @Test
@@ -312,6 +314,32 @@ class ConfigLogicTest {
assertEquals(listOf("s"), out["x-opencode-session"]) assertEquals(listOf("s"), out["x-opencode-session"])
} }
@Test
fun formatHeadersForLogMasksSecretsAndKeepsOthers() {
val line = formatHeadersForLog(
headersOf(
"X-Opencode-Session" to listOf("abc-123"),
"Authorization" to listOf("Bearer super-secret"),
"x-api-key" to listOf("key-1"),
"Cookie" to listOf("session=deadbeef"),
),
)
assertTrue(line.contains("X-Opencode-Session=abc-123"))
assertTrue(line.contains("Authorization=***"))
assertTrue(line.contains("x-api-key=***"))
assertTrue(line.contains("Cookie=***"))
assertFalse(line.contains("super-secret"))
assertFalse(line.contains("deadbeef"))
}
@Test
fun formatHeadersForLogJoinsMultipleValues() {
val line = formatHeadersForLog(
headersOf("X-Custom" to listOf("a", "b")),
)
assertEquals("X-Custom=a|b", line)
}
@Test @Test
fun rebuildFromChunksPreservesAllUpstreamFields() { fun rebuildFromChunksPreservesAllUpstreamFields() {
val sse = """ val sse = """
@@ -354,4 +382,168 @@ class ConfigLogicTest {
val err = out["error"]?.jsonObject ?: error("error block missing") val err = out["error"]?.jsonObject ?: error("error block missing")
assertEquals("Unsupported field 'foo'", err["message"]?.jsonPrimitive?.content) assertEquals("Unsupported field 'foo'", err["message"]?.jsonPrimitive?.content)
} }
@Test
fun sha256MatchesKnownVectors() {
assertEquals(
"e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
sha256Hex(""),
)
assertEquals(
"ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad",
sha256Hex("abc"),
)
assertEquals(
"248d6a61d20638b8e5c026930c3e6039a33ce45964ff2167f6ecedd419db06c1",
sha256Hex("abcdbcdecdefdefgefghfghighijhijkijkljklmklmnlmnomnopnopq"),
)
}
@Test
fun sessionPrefixHashesAreStableWhenHistoryGrows() {
val first = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"S"},{"role":"user","content":"hi"}]}""",
).jsonObject
val second = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"S"},{"role":"user","content":"hi"},{"role":"assistant","content":"yo"},{"role":"user","content":"again"}]}""",
).jsonObject
val h1 = sessionPrefixHashes(first)
val h2 = sessionPrefixHashes(second)
// Цепочка начинается с первого user — system-преамбула не хэшируется.
assertEquals(1, h1.size)
assertEquals(3, h2.size)
assertEquals(h1, h2.take(1))
assertEquals(h1.last(), h2[0])
}
@Test
fun sessionPrefixHashesIgnoreLeadingSystemSoSharedPromptDoesNotCollide() {
val a = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"SAME"},{"role":"user","content":"session A"}]}""",
).jsonObject
val b = Json.parseToJsonElement(
"""{"messages":[{"role":"system","content":"SAME"},{"role":"user","content":"session B"}]}""",
).jsonObject
// разные первые user-сообщения → разные хэши, несмотря на общий system
assertFalse(sessionPrefixHashes(a) == sessionPrefixHashes(b))
}
@Test
fun sessionPrefixHashesFallBackToWholeBodyWithoutMessages() {
val body = Json.parseToJsonElement("""{"model":"m"}""").jsonObject
assertEquals(listOf(sha256Hex(body.toString())), sessionPrefixHashes(body))
}
@Test
fun sessionRegistryReusesIdByLongestCommonPrefix() {
val reg = SessionRegistry()
val first = reg.resolveLocked(listOf("H1"), 0)
assertEquals("H1", first)
// история выросла: [H1, H2, H3] — самый длинный известный префикс H1
val next = reg.resolveLocked(listOf("H1", "H2", "H3"), 1)
assertEquals("H1", next)
// и дальше — id не меняется
val deep = reg.resolveLocked(listOf("H1", "H2", "H3", "H4"), 2)
assertEquals("H1", deep)
}
@Test
fun sessionRegistryCreatesNewIdForDifferentHistory() {
val reg = SessionRegistry()
assertEquals("A1", reg.resolveLocked(listOf("A1"), 0))
assertEquals("B1", reg.resolveLocked(listOf("B1"), 0))
assertEquals("A1", reg.resolveLocked(listOf("A1", "A2"), 0))
}
@Test
fun sessionRegistryEvictsLruWhenOverCapacity() {
val reg = SessionRegistry(maxSessions = 2, ttlMillis = Long.MAX_VALUE)
reg.resolveLocked(listOf("A1"), 0)
reg.resolveLocked(listOf("B1"), 1)
reg.resolveLocked(listOf("C1"), 2) // A вытеснена
// a1 больше неизвестен → новая сессия с id = хэш всей истории A2
assertEquals("A2", reg.resolveLocked(listOf("A1", "A2"), 3))
}
@Test
fun sessionRegistryEvictsExpiredByTtl() {
val reg = SessionRegistry(maxSessions = 100, ttlMillis = 1000)
reg.resolveLocked(listOf("A1"), 0)
reg.resolveLocked(listOf("B1"), 2000) // A протухла
assertEquals("A2", reg.resolveLocked(listOf("A1", "A2"), 2001))
}
@Test
fun parseConfigReadsSessionHeader() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
session_header: x-opencode-session
- id: p2
url: "https://y.ru/api/v1"
models:
- name: m1
upstreams: []
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
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))
}
} }
@@ -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("<think>x</think>" to "", s.feed("<think>x</think>"))
assertEquals("привет\nещё" to "", s.feed("привет\nещё"))
assertEquals("" to "", s.finish())
}
@Test
fun splitWholeBlockMovesInnerTextToReasoning() {
val s = ThinkTagSplitter("split")
assertEquals("a" to "r", s.feed("<think>r</think>a"))
assertEquals("" to "", s.finish())
}
@Test
fun splitTagCutAcrossThreeChunks() {
val s = ThinkTagSplitter("split")
val first = s.feed("привет <t")
assertEquals("привет ", first.first)
assertEquals("", first.second)
assertEquals("" to "", s.feed("hi"))
val third = s.feed("nk>разум")
assertEquals("", third.first)
assertEquals("разум", third.second)
assertEquals("" to "", s.finish())
}
@Test
fun splitUnclosedOpenTagLeavesRestInReasoning() {
val s = ThinkTagSplitter("split")
assertEquals("" to "мысли без конца", s.feed("<think>мысли без конца"))
assertEquals("" to "", s.finish())
// хвост-префикс незакрытого закрывающего тега тоже уезжает в reasoning
val s2 = ThinkTagSplitter("split")
assertEquals("" to "abc", s2.feed("<think>abc</thi"))
assertEquals("" to "</thi", s2.finish())
}
@Test
fun splitHandlesSeveralBlocksInOneResponse() {
val s = ThinkTagSplitter("split")
val first = s.feed("<think>A</think>B<think>C</think>")
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("<think>x</think>y"))
assertEquals("" to "", s.feed("<think>z"))
assertEquals("" to "", s.finish())
}
@Test
fun heldTailThatTurnedOutNotToBeATagFlowsBackAsContent() {
val s = ThinkTagSplitter("split")
// "<thi" — полный префикс открывающего тега: хвост держим
assertEquals("" to "", s.feed("<thi"))
// "x" не продолжает тег: хвост отдаём как обычный текст
assertEquals("<thix" to "", s.feed("x"))
assertEquals("" to "", s.finish())
}
}