4 Commits
2 .. 5

3 changed files with 149 additions and 15 deletions
+4 -4
View File
@@ -12,10 +12,10 @@ repositories {
} }
dependencies { dependencies {
implementation("io.ktor:ktor-server-core:3.1.2") implementation("io.ktor:ktor-server-core:3.5.2")
implementation("io.ktor:ktor-server-cio:3.1.2") implementation("io.ktor:ktor-server-cio:3.5.2")
implementation("io.ktor:ktor-server-content-negotiation:3.1.2") implementation("io.ktor:ktor-server-content-negotiation:3.5.2")
implementation("io.ktor:ktor-serialization-kotlinx-json:3.1.2") implementation("io.ktor:ktor-serialization-kotlinx-json:3.5.2")
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.8.1") implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.8.1")
implementation("ch.qos.logback:logback-classic:1.5.18") implementation("ch.qos.logback:logback-classic:1.5.18")
} }
+8 -1
View File
@@ -3,15 +3,22 @@
# моделей с суффиксом "-no-think", проксирует на RouterAI. Ответ — как есть. # моделей с суффиксом "-no-think", проксирует на RouterAI. Ответ — как есть.
# #
# Запуск на 76.179 (llm-router): podman-compose up -d # Запуск на 76.179 (llm-router): podman-compose up -d
# Образ тянется из нашего реестра (zot): images.binom.pw/llm-proxy:<tag>
services: services:
llm-proxy: llm-proxy:
image: images.binom.pw/llm-proxy:latest image: images.binom.pw/llm-proxy:latest
container_name: llm-proxy container_name: llm-proxy
restart: unless-stopped restart: unless-stopped
network_mode: bifrost_default networks:
- bifrost
environment: environment:
- PORT=8100 - PORT=8100
- UPSTREAM_URL=https://routerai.ru/api/v1 - UPSTREAM_URL=https://routerai.ru/api/v1
- ROUTER_API_KEY=__SET_FROM_CONFIG_DB__ - ROUTER_API_KEY=__SET_FROM_CONFIG_DB__
- EXCLUDED_PROVIDERS=deepseek - EXCLUDED_PROVIDERS=deepseek
- THINKING_MODELS=deepseek/deepseek-v4-flash-0731,deepseek/deepseek-v4-flash - THINKING_MODELS=deepseek/deepseek-v4-flash-0731,deepseek/deepseek-v4-flash
networks:
bifrost:
external: true
name: bifrost_default
+137 -10
View File
@@ -13,11 +13,19 @@ import io.ktor.server.response.respondOutputStream
import io.ktor.server.routing.get import io.ktor.server.routing.get
import io.ktor.server.routing.post import io.ktor.server.routing.post
import io.ktor.server.routing.routing import io.ktor.server.routing.routing
import io.ktor.server.http.HttpRequestLifecycle
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonArray import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive import kotlinx.serialization.json.JsonPrimitive
import kotlinx.serialization.json.jsonArray import kotlinx.serialization.json.jsonArray
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.suspendCancellableCoroutine
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.ensureActive
import kotlin.coroutines.resume
import kotlin.coroutines.resumeWithException
import java.util.concurrent.CompletableFuture
import kotlinx.serialization.json.jsonObject import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive import kotlinx.serialization.json.jsonPrimitive
import java.net.URI import java.net.URI
@@ -56,6 +64,9 @@ fun main() {
.split(",").map { it.trim() }.filter { it.isNotEmpty() }.distinct() .split(",").map { it.trim() }.filter { it.isNotEmpty() }.distinct()
embeddedServer(CIO, port = port, host = "0.0.0.0") { embeddedServer(CIO, port = port, host = "0.0.0.0") {
install(HttpRequestLifecycle) {
cancelCallOnClose = true
}
proxyModule(upstream, apiKey, excluded, thinking) proxyModule(upstream, apiKey, excluded, thinking)
}.start(wait = true) }.start(wait = true)
} }
@@ -103,6 +114,10 @@ private suspend fun handleChat(
return return
} }
val start = System.currentTimeMillis()
val model = try {
json.parseToJsonElement(patched).jsonObject["model"]?.jsonPrimitive?.content ?: "?"
} catch (e: Exception) { "?" }
val req = HttpRequest.newBuilder() val req = HttpRequest.newBuilder()
.uri(URI.create(upstream.trimEnd('/') + "/chat/completions")) .uri(URI.create(upstream.trimEnd('/') + "/chat/completions"))
.header("Authorization", "Bearer $apiKey") .header("Authorization", "Bearer $apiKey")
@@ -111,17 +126,61 @@ private suspend fun handleChat(
.build() .build()
val stream = json.parseToJsonElement(raw).jsonObject["stream"]?.jsonPrimitive?.content == "true" val stream = json.parseToJsonElement(raw).jsonObject["stream"]?.jsonPrimitive?.content == "true"
if (stream) { try {
// SSE-стрим: транслируем как есть, чанк за чанком. if (stream) {
val resp = http.send(req, HttpResponse.BodyHandlers.ofInputStream()) // SSE-стрим: транслируем как есть, чанк за чанком.
val ct = resp.headers().firstValue("content-type").orElse("text/event-stream") val resp = http.send(req, HttpResponse.BodyHandlers.ofInputStream())
call.respondOutputStream(ContentType.parse(ct), HttpStatusCode.fromValue(resp.statusCode())) { val ct = resp.headers().firstValue("content-type").orElse("text/event-stream")
resp.body().use { input -> input.copyTo(this, 8192) } call.respondOutputStream(ContentType.parse(ct), HttpStatusCode.fromValue(resp.statusCode())) {
resp.body().use { input -> input.copyTo(this, 8192) }
}
println("[llm-proxy] chat model=$model status=${resp.statusCode()} в ${System.currentTimeMillis() - start}ms stream=true")
} else {
// Клиент ждёт non-stream ответ, но апстриму шлём stream=true:
// не-стрим генерация у llama.cpp НЕ отменяется обрывом соединения,
// а стрим — отменяется. Так отмена клиента реально рвёт генерацию.
val streamed = json.parseToJsonElement(patched).jsonObject.toMutableMap().apply {
this["stream"] = JsonPrimitive(true)
}
val req2 = HttpRequest.newBuilder()
.uri(URI.create(upstream.trimEnd('/') + "/chat/completions"))
.header("Authorization", "Bearer $apiKey")
.header("Content-Type", "application/json")
.POST(HttpRequest.BodyPublishers.ofString(JsonObject(streamed).toString()))
.build()
val future = http.sendAsync(req2, HttpResponse.BodyHandlers.ofInputStream())
val resp = future.awaitOrCancel()
val body = resp.body()
// Читаем с проверкой отмены: при обрыве клиента ensureActive() бросит
// CancellationException, а finally закроет входной поток — это рвёт
// апстрим-соединение, и llama.cpp отменяет генерацию.
val full = try {
val sb = StringBuilder()
val reader = body.bufferedReader()
while (true) {
currentCoroutineContext().ensureActive()
val line = reader.readLine() ?: break
sb.append(line).append('\n')
}
sb.toString()
} finally {
body.close()
}
val ct = resp.headers().firstValue("content-type").orElse("application/json")
val out = if (ct.contains("text/event-stream")) rebuildFromChunks(full) else full
call.respondBytes(out.toByteArray(), ContentType.parse(ct), HttpStatusCode.fromValue(resp.statusCode()))
println("[llm-proxy] chat model=$model status=${resp.statusCode()} в ${System.currentTimeMillis() - start}ms stream=false")
} }
} else { } catch (e: CancellationException) {
val resp = http.send(req, HttpResponse.BodyHandlers.ofByteArray()) // Клиент оборвал соединение: апстрим-запрос уже отменён через awaitOrCancel.
val ct = resp.headers().firstValue("content-type").orElse("application/json") println("[llm-proxy] chat model=$model ОТМЕНЕНО клиентом в ${System.currentTimeMillis() - start}ms")
call.respondBytes(resp.body(), ContentType.parse(ct), HttpStatusCode.fromValue(resp.statusCode())) throw e
} catch (e: Exception) {
println("[llm-proxy] chat model=$model ОШИБКА: ${e.message} в ${System.currentTimeMillis() - start}ms")
call.respondBytes(
"""{"error":{"message":"upstream: ${e.message}"}}""".toByteArray(),
ContentType.Application.Json, HttpStatusCode.BadGateway,
)
} }
} }
@@ -141,6 +200,12 @@ internal fun patchBody(raw: String, excluded: List<String>, thinking: List<Strin
if (root["reasoning"] == null) { if (root["reasoning"] == null) {
root["reasoning"] = JsonObject(mapOf("enabled" to JsonPrimitive(false))) root["reasoning"] = JsonObject(mapOf("enabled" to JsonPrimitive(false)))
} }
// SGLANG_COMPAT: Qwen3.8 глушится только через chat_template_kwargs.enable_thinking=false
if (System.getenv("SGLANG_COMPAT") == "true") {
val ctk = root["chat_template_kwargs"]?.jsonObject?.toMutableMap() ?: mutableMapOf()
ctk["enable_thinking"] = JsonPrimitive(false)
root["chat_template_kwargs"] = JsonObject(ctk)
}
} }
val provider = root["provider"]?.jsonObject?.toMutableMap() ?: mutableMapOf() val provider = root["provider"]?.jsonObject?.toMutableMap() ?: mutableMapOf()
@@ -194,3 +259,65 @@ internal fun patchModelsCatalog(raw: String, thinking: List<String>): String {
root["data"] = JsonArray(data + copies) root["data"] = JsonArray(data + copies)
return JsonObject(root).toString() return JsonObject(root).toString()
} }
/**
* Ожидание CompletableFuture с пробросом отмены корутины на апстрим-запрос:
* если клиент оборвал соединение (Ktor отменяет корутину), рвём и апстрим —
* upstream (llama.cpp/sglang) видит обрыв и отменяет генерацию (слот свободен).
*/
private suspend fun <T> CompletableFuture<T>.awaitOrCancel(): T =
suspendCancellableCoroutine { cont ->
this.whenComplete { res, err ->
if (err != null) cont.resumeWithException(err) else cont.resume(res)
}
cont.invokeOnCancellation { this.cancel(true) }
}
/**
* Собрать полный chat.completion из SSE-чанков апстрима (для non-stream клиентов).
*/
private fun rebuildFromChunks(sse: String): String {
var content = StringBuilder()
var reasoning = StringBuilder()
var finish = "stop"
var id = ""
var model = ""
val created = System.currentTimeMillis() / 1000
sse.lineSequence().forEach { line ->
if (!line.startsWith("data:")) return@forEach
val data = line.removePrefix("data:").trim()
if (data.isEmpty() || data == "[DONE]") return@forEach
try {
val obj = json.parseToJsonElement(data).jsonObject
if (id.isEmpty()) id = obj["id"]?.jsonPrimitive?.content ?: ""
if (model.isEmpty()) model = obj["model"]?.jsonPrimitive?.content ?: ""
val choice = obj["choices"]?.jsonArray?.firstOrNull()?.jsonObject
if (choice != null) {
choice["finish_reason"]?.jsonPrimitive?.content
?.takeIf { it.isNotEmpty() && it != "null" }?.let { finish = it }
val delta = choice["delta"]?.jsonObject
delta?.get("content")?.jsonPrimitive?.content
?.takeIf { it != "null" }?.let { content.append(it) }
delta?.get("reasoning_content")?.jsonPrimitive?.content
?.takeIf { it != "null" }?.let { reasoning.append(it) }
}
} catch (_: Exception) {}
}
val msg = JsonObject(mutableMapOf(
"role" to JsonPrimitive("assistant"),
"content" to JsonPrimitive(content.toString()),
"reasoning" to JsonPrimitive(reasoning.toString()),
))
val choice = JsonObject(mutableMapOf(
"index" to JsonPrimitive(0),
"message" to msg,
"finish_reason" to JsonPrimitive(finish),
))
return JsonObject(mutableMapOf(
"id" to JsonPrimitive(id),
"object" to JsonPrimitive("chat.completion"),
"created" to JsonPrimitive(created),
"model" to JsonPrimitive(model),
"choices" to JsonArray(listOf(choice)),
)).toString()
}