diff --git a/README.md b/README.md index 1a8fcfe..03eb4f1 100644 --- a/README.md +++ b/README.md @@ -8,19 +8,33 @@ ## Конфиг Конфигурация — в YAML-файле (по умолчанию `config.yaml` в каталоге проекта, -CWD; переопределяется env `CONFIG_PATH`). Три блока: `providers`, `upstreams`, +CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` (интерфейс биндинга +`host` + `port`, дефолты `0.0.0.0` / `8100`), `providers`, `upstreams`, `models`. Старый механизм env-переменных (`UPSTREAM_URL`, `ROUTER_API_KEY`, -`EXCLUDED_PROVIDERS`, `THINKING_MODELS`) удалён — его поведение теперь в -декларативном `patch` и блоке `upstreams`. +`EXCLUDED_PROVIDERS`, `THINKING_MODELS`, `PORT`) удалён — его поведение теперь +в декларативном `patch` и блоке `upstreams`, а порт/интерфейс — в блоке +`server`. + +Лимит конкурентности `max_concurrency` можно задать и на провайдере (лимит по +умолчанию для его апстримов), и на апстриме (перекрывает провайдерский); без +обоих — безлимит. | Переменная | Default | Описание | |---|---|---| -| `PORT` | 8100 | Порт сервера | | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | +| Ключ конфига | Default | Описание | +|---|---|---| +| `server.host` | `0.0.0.0` | Интерфейс/адрес биндинга | +| `server.port` | `8100` | Порт сервера | + +Таймаут на запрос к нейронке (весь ответ, включая стриминг) — константа +`UPSTREAM_REQUEST_TIMEOUT_MS` = 5 минут; без неё дефолт CIO-движка Ktor — 15 секунд, +и длинные генерации обрываются. Значение логируется при старте. + ## Сборка и деплой -- CI (Gitea Actions, на release): `./gradlew shadowJar` → образ +- CI (Gitea Actions, на release): `./gradlew fatJar` → образ `images.binom.pw/llm-proxy:` (zot). - Запуск на 76.179 (llm-router): `podman-compose up -d` из `podman-compose.yaml` (образ тянется из `images.binom.pw`, сеть `bifrost_default` — туда же Bifrost diff --git a/podman-compose.yaml b/podman-compose.yaml index 34e7a25..e9ac146 100644 --- a/podman-compose.yaml +++ b/podman-compose.yaml @@ -12,7 +12,6 @@ services: networks: - bifrost environment: - - PORT=8100 # Путь к конфигу внутри контейнера (файл монтируется ниже). - CONFIG_PATH=/config/config.yaml volumes: diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt index ca1c707..08ff803 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt @@ -54,11 +54,11 @@ import kotlin.time.TimeSource * вмердживаются `patch` (provider → upstream → model). * * Env (только эти): - * - PORT (default 8100) * - CONFIG_PATH путь к YAML; default `config.yaml` в каталоге проекта (CWD) + * + * Порт/интерфейс биндинга задаются блоком `server` в YAML (см. CONFIG.md). */ fun main() { - val port = getEnv("PORT")?.toIntOrNull() ?: 8100 val path = getEnv("CONFIG_PATH") ?: "config.yaml" val root = Yaml.decodeYamlFromString(readConfigText(path)) val config = parseConfig(root) @@ -78,15 +78,23 @@ fun main() { } } - val active = config.upstreams.associate { up -> up.id to UpstreamCounter(up.max_concurrency ?: Int.MAX_VALUE) } + val active = config.upstreams.associate { up -> + up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider])) + } log.info { "[llm-proxy] загружено: providers=${config.providers.size}, " + "upstreams=${config.upstreams.size}, models=${config.models.size} (config=$path)" } + log.info { + "[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}" + } + log.info { + "[llm-proxy] upstream: request_timeout=${UPSTREAM_REQUEST_TIMEOUT_MS}ms" + } val http = createHttpClient() - startServer(port) { + startServer(config.server.host, config.server.port) { proxyModule(config, providersById, upstreamsById, active, http) } } @@ -308,12 +316,16 @@ internal fun merge(base: JsonObject, patch: JsonObject): JsonObject { internal fun resolveEnv(s: String): String = """\$\{([^}]+)\}""".toRegex().replace(s) { m -> getEnv(m.groupValues[1]) ?: "" } -/** Атомарно занять слот у апстрима (по max_concurrency); false, если все заняты. */ -internal fun tryClaim(up: UpstreamConf, active: Map): Boolean { - val counter = active.getValue(up.id) - val limit = up.max_concurrency ?: Int.MAX_VALUE - return counter.tryClaim(limit) -} +/** Атомарно занять слот у апстрима (по эффективному лимиту); false, если все заняты. */ +internal fun tryClaim(up: UpstreamConf, active: Map): Boolean = + active.getValue(up.id).tryClaim() + +/** + * Эффективный лимит конкурентности апстрима: значение у апстрима (модели), если + * задано; иначе у провайдера; иначе безлимит. + */ +internal fun effectiveConcurrencyLimit(up: UpstreamConf, provider: ProviderConf?): Int = + up.max_concurrency ?: provider?.max_concurrency ?: Int.MAX_VALUE /** Освободить слот апстрима (в finally по завершении проксирования). */ internal fun release(up: UpstreamConf, active: Map) { @@ -325,10 +337,10 @@ class UpstreamCounter(private val limit: Int, initial: Int = 0) { @Volatile private var count: Int = initial - fun tryClaim(max: Int): Boolean { + fun tryClaim(): Boolean { if (!lock.tryLock()) return false try { - return if (count >= max) false else { + return if (count >= limit) false else { count++ true } @@ -483,6 +495,7 @@ data class ProviderConf( val id: String, val url: String, val key: String = "", + val max_concurrency: Int? = null, val patch: JsonObject? = null, ) @@ -500,7 +513,13 @@ data class ModelConf( val patch: JsonObject? = null, ) +data class ServerConf( + val host: String = "0.0.0.0", + val port: Int = 8100, +) + data class Config( + val server: ServerConf, val providers: List, val upstreams: List, val models: List, @@ -521,12 +540,23 @@ internal fun parseConfig(root: YamlElement): Config { return (v as? YamlList)?.map { it } ?: emptyList() } + val serverMap = (top["server"] as? YamlMap)?.toMap() + val server = if (serverMap != null) { + ServerConf( + host = serverMap.strOrNull("host") ?: "0.0.0.0", + port = serverMap.strOrNull("port")?.toIntOrNull() ?: 8100, + ) + } else { + ServerConf() + } + val providers = list("providers").map { entry -> val m = (entry as YamlMap).toMap() ProviderConf( id = m.str("id"), url = m.str("url"), key = m.strOrNull("key") ?: "", + max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, ) } @@ -551,7 +581,7 @@ internal fun parseConfig(root: YamlElement): Config { ) } - return Config(providers, upstreams, models) + return Config(server, providers, upstreams, models) } /** YamlMap -> Map (ключи YAML — строковые скаляры). */ diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Platform.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Platform.kt index 51d29ef..582d32e 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Platform.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Platform.kt @@ -7,4 +7,4 @@ import io.ktor.server.application.Application expect fun createHttpClient(): HttpClient /** Платформенный запуск Ktor-сервера (движок задаётся в actual). */ -expect fun startServer(port: Int, module: Application.() -> Unit) +expect fun startServer(host: String, port: Int, module: Application.() -> Unit) diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt index 61c0cc1..7e5c593 100644 --- a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt +++ b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt @@ -35,6 +35,7 @@ class ConfigLogicTest { - id: p1 url: "https://x.ru/api/v1" key: "k" + max_concurrency: 4 patch: provider: allow_fallbacks: false @@ -61,6 +62,7 @@ class ConfigLogicTest { assertEquals(1, cfg.providers.size) assertEquals("https://x.ru/api/v1", cfg.providers[0].url) + assertEquals(4, cfg.providers[0].max_concurrency) assertEquals(2, cfg.upstreams.size) assertEquals(2, cfg.upstreams[0].max_concurrency) assertEquals(listOf("u1", "bad"), cfg.models[0].upstreams) @@ -70,6 +72,53 @@ class ConfigLogicTest { assertEquals("""{"reasoning":{"enabled":false}}""", cfg.models[0].patch.toString()) } + @Test + fun parseConfigReadsServerBlock() { + val yaml = """ + server: + host: 127.0.0.1 + port: 9200 + models: + - name: m1 + upstreams: [] + """.trimIndent() + + val cfg = parseConfig(Yaml.decodeYamlFromString(yaml)) + + assertEquals("127.0.0.1", cfg.server.host) + assertEquals(9200, cfg.server.port) + } + + @Test + fun parseConfigServerDefaultsWhenBlockAbsent() { + val yaml = """ + models: + - name: m1 + upstreams: [] + """.trimIndent() + + val cfg = parseConfig(Yaml.decodeYamlFromString(yaml)) + + assertEquals("0.0.0.0", cfg.server.host) + assertEquals(8100, cfg.server.port) + } + + @Test + fun parseConfigServerFieldDefaultsPerField() { + val yaml = """ + server: + host: 192.168.88.10 + models: + - name: m1 + upstreams: [] + """.trimIndent() + + val cfg = parseConfig(Yaml.decodeYamlFromString(yaml)) + + assertEquals("192.168.88.10", cfg.server.host) + assertEquals(8100, cfg.server.port) + } + @Test fun parseConfigOmitsPatchWhenAbsent() { val yaml = """ @@ -93,7 +142,7 @@ class ConfigLogicTest { @Test fun buildBodySubstitutesModelAndAppliesLayersInOrder() { val provider = ProviderConf( - "p1", "https://x", "", + "p1", "https://x", "", null, Json.parseToJsonElement("""{"provider":{"allow_fallbacks":false}}""").jsonObject, ) val up = UpstreamConf( @@ -133,7 +182,7 @@ class ConfigLogicTest { @Test fun tryClaimRespectsMaxConcurrencyAndReleaseFreesSlot() { - val active = mapOf("u1" to UpstreamCounter(0)) + val active = mapOf("u1" to UpstreamCounter(1)) val up = UpstreamConf("u1", "p", "m", 1, null) assertTrue(tryClaim(up, active)) assertFalse(tryClaim(up, active)) @@ -143,15 +192,37 @@ class ConfigLogicTest { @Test fun tryClaimUnlimitedWhenMaxConcurrencyIsNull() { - val active = mapOf("u2" to UpstreamCounter(0)) + val active = mapOf("u2" to UpstreamCounter(Int.MAX_VALUE)) val up = UpstreamConf("u2", "p", "m", null, null) assertTrue(tryClaim(up, active)) assertTrue(tryClaim(up, active)) } + @Test + fun effectiveConcurrencyLimitPrefersUpstreamThenProvider() { + val providerWithLimit = ProviderConf("p1", "https://x", "", 3, null) + val providerWithoutLimit = ProviderConf("p2", "https://x", "", null, null) + assertEquals( + 2, + effectiveConcurrencyLimit(UpstreamConf("u1", "p1", "m", 2, null), providerWithLimit), + ) + assertEquals( + 3, + effectiveConcurrencyLimit(UpstreamConf("u1", "p1", "m", null, null), providerWithLimit), + ) + assertEquals( + Int.MAX_VALUE, + effectiveConcurrencyLimit(UpstreamConf("u1", "p2", "m", null, null), providerWithoutLimit), + ) + assertEquals( + Int.MAX_VALUE, + effectiveConcurrencyLimit(UpstreamConf("u1", "missing", "m", null, null), null), + ) + } + @Test fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() { - val active = mapOf("u1" to UpstreamCounter(0), "u2" to UpstreamCounter(0)) + val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val pool = listOf( UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 2, null), @@ -164,7 +235,7 @@ class ConfigLogicTest { @Test fun pickFreeUpstreamSkipsExcluded() { - val active = mapOf("u1" to UpstreamCounter(0), "u2" to UpstreamCounter(0)) + val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val pool = listOf( UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 2, null), @@ -183,7 +254,7 @@ class ConfigLogicTest { @Test fun pickFreeUpstreamImplementsFailoverOrder() { // dead исключён (упал ранее) — выбирается следующий живой u1 - val active = mapOf("dead" to UpstreamCounter(0), "u1" to UpstreamCounter(0)) + val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1)) val pool = listOf( UpstreamConf("dead", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null), @@ -194,7 +265,7 @@ class ConfigLogicTest { @Test fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() { - val active = mapOf("u1" to UpstreamCounter(0), "u2" to UpstreamCounter(0)) + val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1)) val pool = listOf( UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 1, null), diff --git a/src/jvmMain/kotlin/pw/binom/llmproxy/Platform.kt b/src/jvmMain/kotlin/pw/binom/llmproxy/Platform.kt index e2091fc..51a330f 100644 --- a/src/jvmMain/kotlin/pw/binom/llmproxy/Platform.kt +++ b/src/jvmMain/kotlin/pw/binom/llmproxy/Platform.kt @@ -6,10 +6,14 @@ import io.ktor.server.application.Application import io.ktor.server.cio.CIO as ServerCIO import io.ktor.server.engine.embeddedServer -actual fun createHttpClient(): HttpClient = HttpClient(CIO) +actual fun createHttpClient(): HttpClient = HttpClient(CIO) { + engine { + requestTimeout = UPSTREAM_REQUEST_TIMEOUT_MS + } +} -actual fun startServer(port: Int, module: Application.() -> Unit) { - embeddedServer(ServerCIO, port = port, host = "0.0.0.0") { +actual fun startServer(host: String, port: Int, module: Application.() -> Unit) { + embeddedServer(ServerCIO, port = port, host = host) { module() }.start(wait = true) } diff --git a/src/linuxX64Main/kotlin/pw/binom/llmproxy/Platform.kt b/src/linuxX64Main/kotlin/pw/binom/llmproxy/Platform.kt index e2091fc..51a330f 100644 --- a/src/linuxX64Main/kotlin/pw/binom/llmproxy/Platform.kt +++ b/src/linuxX64Main/kotlin/pw/binom/llmproxy/Platform.kt @@ -6,10 +6,14 @@ import io.ktor.server.application.Application import io.ktor.server.cio.CIO as ServerCIO import io.ktor.server.engine.embeddedServer -actual fun createHttpClient(): HttpClient = HttpClient(CIO) +actual fun createHttpClient(): HttpClient = HttpClient(CIO) { + engine { + requestTimeout = UPSTREAM_REQUEST_TIMEOUT_MS + } +} -actual fun startServer(port: Int, module: Application.() -> Unit) { - embeddedServer(ServerCIO, port = port, host = "0.0.0.0") { +actual fun startServer(host: String, port: Int, module: Application.() -> Unit) { + embeddedServer(ServerCIO, port = port, host = host) { module() }.start(wait = true) }