Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e251ca5326 | |||
| 0fb99060f9 |
@@ -5,3 +5,9 @@ build/
|
|||||||
# локальный конфиг с секретами (ключ API) — не коммитим
|
# локальный конфиг с секретами (ключ API) — не коммитим
|
||||||
config.yaml
|
config.yaml
|
||||||
|
|
||||||
|
# IDE / tooling
|
||||||
|
.idea/
|
||||||
|
.kotlin/
|
||||||
|
.cortexkit/
|
||||||
|
.veai/
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,391 @@
|
|||||||
|
# Конфигурация llm-proxy (переход на YAML)
|
||||||
|
|
||||||
|
## Идея
|
||||||
|
|
||||||
|
Вместо переменных среды конфигурация задаётся YAML-файлом. В нём три блока:
|
||||||
|
|
||||||
|
1. **`providers`** — объявлены бэкенды (url + ключ) с нашим внутренним
|
||||||
|
человекочитаемым `id`. По этому `id` ссылаются апстримы. Опционально —
|
||||||
|
`max_concurrency`: лимит по умолчанию для апстримов этого провайдера.
|
||||||
|
2. **`upstreams`** — каталог «удалённых» моделей. Каждая запись: внутренний `id`,
|
||||||
|
ссылка на `providers[].id`, реальное имя модели у провайдера (`model`) и
|
||||||
|
опциональный лимит конкурентности `max_concurrency`. Это и есть то, на что
|
||||||
|
мы проксируем.
|
||||||
|
3. **`models`** — модели, видимые клиенту. Для каждой — «витринное» имя (`name`),
|
||||||
|
список `upstreams` (ссылки на `upstreams[].id`, можно несколько).
|
||||||
|
|
||||||
|
4. **`patch`** — не отдельный блок, а опциональное поле, доступное **на всех
|
||||||
|
уровнях** (`providers`, `upstreams`, `models`). Это JSON-подобная структура
|
||||||
|
(`kotlinx.serialization.json.JsonObject`, в YAML пишется как обычный мап),
|
||||||
|
которая **аккуратно вмердживается** в тело запроса. Пишется нативным YAML,
|
||||||
|
без вложенного JSON-в-строке.
|
||||||
|
|
||||||
|
При запросе к модели прокси берёт **первый из её `upstreams`, на котором прямо
|
||||||
|
сейчас есть свободный слот** (по `max_concurrency`). Это позволяет держать
|
||||||
|
провайдеров, допускающих лишь один одновременный запрос (`max_concurrency: 1`),
|
||||||
|
рядом с многопоточными — и не превышать их лимиты.
|
||||||
|
|
||||||
|
**Лимит конкурентности** апстрима берётся в порядке убывания специфичности:
|
||||||
|
`upstreams[].max_concurrency` (самый конкретный) → `providers[].max_concurrency`
|
||||||
|
(провайдера, на который ссылается апстрим) → безлимит. Т.е. задан у модели —
|
||||||
|
берём её; у модели нет, но есть у провайдера — берём провайдерский.
|
||||||
|
|
||||||
|
**Сборка тела запроса** — берём исходный JSON клиента и накладываем `patch`
|
||||||
|
**послойно, в порядке возрастания специфичности**: сначала `provider.patch`
|
||||||
|
(всё, что идёт через этого провайдера), затем `upstream.patch` (конкретная
|
||||||
|
апстрим-модель),最后 `model.patch` (самый приоритетный, накладывается последним).
|
||||||
|
Мерж — глубокий: вложенные объекты сливаются рекурсивно, скаляры/массивы по ключу
|
||||||
|
заменяются значением из `patch`.
|
||||||
|
|
||||||
|
**Маршрутизация целиком опирается на внутренний учёт** — мы сами считаем, сколько
|
||||||
|
запросов сейчас идёт на каждый апстрим, и по этим счётчикам решаем, брать запрос
|
||||||
|
или нет. Мы **не полагаемся** на `503`/`429` от самого апстрима, чтобы узнать,
|
||||||
|
что он перегружен: такой ответ от провайдера — это уже сбой, а не способ
|
||||||
|
управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`, мы
|
||||||
|
освобождаем его слот и, при наличии, пробуем **следующий свободный** из списка
|
||||||
|
модели (фейловер); если свободных не осталось — отдаём `503` от прокси.
|
||||||
|
|
||||||
|
## Формат файла
|
||||||
|
|
||||||
|
Путь к файлу: по умолчанию ищем `config.yaml` в **каталоге проекта** (текущем
|
||||||
|
рабочем каталоге, откуда запущен процесс). Переопределить можно через env
|
||||||
|
`CONFIG_PATH` — тогда берётся указанный путь (абсолютный или относительный).
|
||||||
|
|
||||||
|
> **URL — в кавычках.** В YAML значение вида `https://...` парсится как вложенный
|
||||||
|
> мап (`https:` — ключ), поэтому `url` (и любое значение с `://`) **обязательно
|
||||||
|
> берите в кавычки**: `url: "https://routerai.ru/api/v1"`.
|
||||||
|
|
||||||
|
### Блок `server` — биндинг HTTP-сервера
|
||||||
|
|
||||||
|
Блок `server` задаёт, **на какой интерфейс (`host`) и порт (`port`)** слушает
|
||||||
|
прокси. Оба поля опциональны: дефолт `host: 0.0.0.0`, `port: 8100`. При старте
|
||||||
|
сервер логирует фактический биндинг:
|
||||||
|
`[llm-proxy] server: bind host=<host> port=<port>`.
|
||||||
|
|
||||||
|
```yaml
|
||||||
|
server:
|
||||||
|
host: 0.0.0.0 # интерфейс/адрес биндинга (дефолт 0.0.0.0)
|
||||||
|
port: 8100 # порт (дефолт 8100)
|
||||||
|
|
||||||
|
# 1) Бэкенды. id — наш внутренний идентификатор (на него ссылаются апстримы).
|
||||||
|
providers:
|
||||||
|
- id: routerai # любая строка, уникальная в рамках файла
|
||||||
|
url: "https://routerai.ru/api/v1"
|
||||||
|
key: "sk-..." # Bearer-ключ; можно подставлять из env
|
||||||
|
max_concurrency: 4 # опционально; лимит по умолчанию для апстримов
|
||||||
|
patch: # уровень провайдера: ко всем его запросам
|
||||||
|
provider:
|
||||||
|
allow_fallbacks: false
|
||||||
|
|
||||||
|
- id: local-llama
|
||||||
|
url: "http://10.0.0.5:8080/v1"
|
||||||
|
key: "" # пусто, если бэкенд без авторизации
|
||||||
|
max_concurrency: 1 # и локалка — один слот за раз
|
||||||
|
|
||||||
|
# 2) Каталог апстрим-моделей. id — наш внутренний id (на него ссылаются модели).
|
||||||
|
upstreams:
|
||||||
|
- id: routerai-gpt4o # внутренний id апстрима
|
||||||
|
provider: routerai # ссылка на providers[].id
|
||||||
|
model: gpt-4o # реальное имя модели у провайдера
|
||||||
|
max_concurrency: 4 # опционально; сколько одновременных запросов
|
||||||
|
# допустимо (null/0 = безлимит)
|
||||||
|
patch: # уровень апстрима
|
||||||
|
provider:
|
||||||
|
ignore: [deepseek]
|
||||||
|
|
||||||
|
- id: local-qwen
|
||||||
|
provider: local-llama
|
||||||
|
model: qwen2.5-72b-instruct
|
||||||
|
max_concurrency: 1 # однослотовый провайдер: 1 запрос за раз
|
||||||
|
|
||||||
|
# 3) Модели, видимые клиенту.
|
||||||
|
models:
|
||||||
|
- name: my-gpt # имя, под которым клиент запрашивает модель
|
||||||
|
upstreams: # порядок = приоритет: сначала локальная видюха,
|
||||||
|
- local-qwen # потом (фоллбэк) внешний провайдер
|
||||||
|
- routerai-gpt4o
|
||||||
|
patch: # уровень модели (накладывается последним)
|
||||||
|
reasoning:
|
||||||
|
enabled: false
|
||||||
|
|
||||||
|
- name: deepseek-fast-no-think
|
||||||
|
upstreams:
|
||||||
|
- routerai-gpt4o
|
||||||
|
patch:
|
||||||
|
reasoning:
|
||||||
|
enabled: false
|
||||||
|
|
||||||
|
- name: local-qwen-only # только локалка, без фоллбэка
|
||||||
|
upstreams:
|
||||||
|
- local-qwen
|
||||||
|
# patch необязателен на любом уровне — можно не указывать
|
||||||
|
```
|
||||||
|
|
||||||
|
> **Приоритет/фоллбэк.** Порядок записей в `upstreams` модели — это приоритет:
|
||||||
|
> прокси берёт **первый свободный по порядку**. Типовой сценарий — сначала своя
|
||||||
|
> локальная видюха (`max_concurrency: 1`, занята → следующий), а внешний провайдер
|
||||||
|
> идёт вторым и принимает запрос, только когда локалка занята (или упала).
|
||||||
|
> Чтобы внешний НЕ использовался, пока локалка свободна, — просто ставь локалку
|
||||||
|
> первой; фоллбэк сработает автоматически при `claim() == null` у локалки.
|
||||||
|
|
||||||
|
`patch` опционален на **любом** уровне (`providers` / `upstreams` / `models`):
|
||||||
|
если ни одного нет — запрос проксируется как есть (исходное тело клиента).
|
||||||
|
|
||||||
|
### Пример сборки тела (многослойный `patch`)
|
||||||
|
|
||||||
|
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
|
||||||
|
`routerai-gpt4o` (провайдер `routerai`). Клиент шлёт:
|
||||||
|
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"model": "my-gpt",
|
||||||
|
"messages": [{ "role": "user", "content": "привет" }],
|
||||||
|
"temperature": 0.7
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
Слои `patch`, накладываемые в порядке `provider → upstream → model`:
|
||||||
|
|
||||||
|
| Слой | Что добавляет |
|
||||||
|
|---|---|
|
||||||
|
| `providers.routerai.patch` | `provider.allow_fallbacks = false` |
|
||||||
|
| `upstreams.routerai-gpt4o.patch` | `provider.ignore = ["deepseek"]` |
|
||||||
|
| `models.my-gpt.patch` | `reasoning.enabled = false` |
|
||||||
|
|
||||||
|
Плюс подмена `model: "my-gpt"` → `model: "gpt-4o"` (реальное имя апстрима).
|
||||||
|
Глубокий мерж сливает `provider` из двух слоёв в один объект. **Итоговое тело**,
|
||||||
|
уходящее на `https://routerai.ru/api/v1/chat/completions`:
|
||||||
|
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"model": "gpt-4o",
|
||||||
|
"messages": [{ "role": "user", "content": "привет" }],
|
||||||
|
"temperature": 0.7,
|
||||||
|
"provider": {
|
||||||
|
"allow_fallbacks": false,
|
||||||
|
"ignore": ["deepseek"]
|
||||||
|
},
|
||||||
|
"reasoning": {
|
||||||
|
"enabled": false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
Если бы вместо `routerai-gpt4o` сработал фоллбэк на `local-qwen` — слой
|
||||||
|
`upstreams.routerai-gpt4o.patch` не применился бы (нет у локалки), и в теле не
|
||||||
|
было бы `provider.ignore`/`allow_fallbacks`; модель стала бы `qwen2.5-72b-instruct`,
|
||||||
|
а `model.patch` (`reasoning.enabled=false`) остался бы — он не зависит от апстрима.
|
||||||
|
|
||||||
|
## Ссылка на env в ключах
|
||||||
|
|
||||||
|
Чтобы не хранить ключи в файле, `key` можно резолвить из переменной среды
|
||||||
|
(подстановка вида `${ROUTER_API_KEY}`):
|
||||||
|
|
||||||
|
```yaml
|
||||||
|
providers:
|
||||||
|
- id: routerai
|
||||||
|
url: "https://routerai.ru/api/v1"
|
||||||
|
key: "${ROUTER_API_KEY}"
|
||||||
|
```
|
||||||
|
|
||||||
|
## Что меняется в коде (`Main.kt`)
|
||||||
|
|
||||||
|
> **Старая концепция выпиливается целиком.** Текущий код — это «один апстрим на
|
||||||
|
> одну переменную среды»: `UPSTREAM_URL`, `ROUTER_API_KEY`, `EXCLUDED_PROVIDERS`,
|
||||||
|
> `THINKING_MODELS`, функция `patchBody(...)`, чтение `SGLANG_COMPAT`, а также
|
||||||
|
> дублирование каталога `/v1/models` с суффиксом `-no-think`. Всё это **удаляется**
|
||||||
|
> без обратной совместимости — сервис становится YAML-декларативным роутером
|
||||||
|
> (разделы ниже). Env-переменные конфигурации провайдеров/моделей больше не
|
||||||
|
> поддерживаются; остаётся только `CONFIG_PATH` (порт/интерфейс — в блоке
|
||||||
|
> `server` файла).
|
||||||
|
|
||||||
|
### 1. Модель конфига (сериализуемые классы)
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
// patch — универсальная JSON-подобная структура (kotlinx JsonObject),
|
||||||
|
// декодируется из YAML-мапа. null = патч отсутствует.
|
||||||
|
@Serializable
|
||||||
|
data class ProviderConf(
|
||||||
|
val id: String,
|
||||||
|
val url: String,
|
||||||
|
val key: String = "",
|
||||||
|
val max_concurrency: Int? = null, // лимит по умолчанию для апстримов провайдера
|
||||||
|
val patch: JsonObject? = null, // ко всем запросам провайдера
|
||||||
|
)
|
||||||
|
|
||||||
|
@Serializable
|
||||||
|
data class UpstreamConf(
|
||||||
|
val id: String, // наш внутренний id апстрима
|
||||||
|
val provider: String, // ссылка на ProviderConf.id
|
||||||
|
val model: String, // реальное имя модели у провайдера
|
||||||
|
val max_concurrency: Int? = null,// опционально; null/0 = безлимит
|
||||||
|
val patch: JsonObject? = null, // к запросам этой апстрим-модели
|
||||||
|
)
|
||||||
|
|
||||||
|
@Serializable
|
||||||
|
data class ModelConf(
|
||||||
|
val name: String, // «витринное» имя для клиента
|
||||||
|
val upstreams: List<String>, // ссылки на UpstreamConf.id
|
||||||
|
val patch: JsonObject? = null, // к запросам этой модели (самый приоритетный)
|
||||||
|
)
|
||||||
|
|
||||||
|
@Serializable
|
||||||
|
data class Config(
|
||||||
|
val providers: List<ProviderConf> = emptyList(),
|
||||||
|
val upstreams: List<UpstreamConf> = emptyList(),
|
||||||
|
val models: List<ModelConf> = emptyList(),
|
||||||
|
)
|
||||||
|
```
|
||||||
|
|
||||||
|
> **Тип `patch`.** В YAML пишется как обычный мап, а yamlkt декодирует его в
|
||||||
|
> `kotlinx.serialization.json.JsonObject` (через `JsonObject.serializer()` /
|
||||||
|
> мост yamlkt→JsonElement). Так `patch` — это аккуратная JSON-структура, а не
|
||||||
|
> строка, и мержить её в тело запроса тривиально.
|
||||||
|
|
||||||
|
Глубокий мерж двух `JsonObject` (поля `patch` перезаписывают/дополняют `base`
|
||||||
|
рекурсивно по вложенным объектам):
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
fun merge(base: JsonObject, patch: JsonObject): JsonObject {
|
||||||
|
val merged = base.toMutableMap()
|
||||||
|
for ((k, v) in patch) {
|
||||||
|
merged[k] = when {
|
||||||
|
v is JsonObject && merged[k] is JsonObject ->
|
||||||
|
merge(merged[k] as JsonObject, v) // рекурсивно для вложенных
|
||||||
|
else -> v // скаляр/массив — заменяем
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return JsonObject(merged)
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
Загрузка:
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
// Дефолт — config.yaml в каталоге проекта (CWD). CONFIG_PATH переопределяет путь.
|
||||||
|
val path = System.getenv("CONFIG_PATH") ?: "config.yaml"
|
||||||
|
val configFile = File(path) // относительный путь резолвится от CWD проекта
|
||||||
|
if (!configFile.exists()) error("config not found: ${configFile.absolutePath}")
|
||||||
|
// Весь YAML читается как YamlElement-дерево, затем маппится в Config вручную
|
||||||
|
// (поле patch сразу конвертируется в kotlinx JsonObject). Вложенный
|
||||||
|
// @Serializable-класс с полем YamlElement yamlkt читает некорректно.
|
||||||
|
val root = Yaml.decodeYamlFromString(configFile.readText(Charsets.UTF_8))
|
||||||
|
val config = parseConfig(root)
|
||||||
|
val providersById = config.providers.associateBy { it.id }
|
||||||
|
val upstreamsById = config.upstreams.associateBy { it.id }
|
||||||
|
```
|
||||||
|
|
||||||
|
### 2. Учёт конкурентности (runtime)
|
||||||
|
|
||||||
|
На старте заводим счётчик активных запросов на каждый апстрим — это **единственный
|
||||||
|
источник истины** про занятость. Никаких опросов апстрима и реакции на его `503`.
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
// Эффективный лимит апстрима резолвится при старте: свой max_concurrency, иначе
|
||||||
|
// провайдерский, иначе безлимит. На счётчике лежит уже итоговый лимит.
|
||||||
|
val active = config.upstreams.associate {
|
||||||
|
it.id to UpstreamCounter(
|
||||||
|
it.max_concurrency
|
||||||
|
?: providersById[it.provider]?.max_concurrency
|
||||||
|
?: Int.MAX_VALUE,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
Отбор свободного апстрима — атомарно: пытаемся инкрементировать счётчик, только
|
||||||
|
если он меньше лимита. Если ни один из `upstreams` модели не свободен — отдаём
|
||||||
|
`503` (`all upstreams busy`), слоты при этом не трогаем. Слот освобождается
|
||||||
|
(`decrement`) в `finally` по завершении проксирования (успех/ошибка/отмена).
|
||||||
|
|
||||||
|
```kotlin
|
||||||
|
// Атомарно занимает слот у первого свободного апстрима; вернёт null, если все заняты.
|
||||||
|
fun claim(up: List<UpstreamConf>): UpstreamConf? = up.firstOrNull { u ->
|
||||||
|
val counter = active.getValue(u.id)
|
||||||
|
val cur = counter.get()
|
||||||
|
cur < counter.limit && counter.compareAndSet(cur, cur + 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Освободить слот (в finally).
|
||||||
|
fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
|
||||||
|
```
|
||||||
|
|
||||||
|
### 3. `handleChat` — выбор апстрима, подмена модели и мерж `patch`
|
||||||
|
|
||||||
|
Вместо текущей логики с `EXCLUDED_PROVIDERS` / `THINKING_MODELS`:
|
||||||
|
|
||||||
|
- берём из тела `model`;
|
||||||
|
- ищем `config.models.first { it.name == model }` (404, если нет);
|
||||||
|
- резолвим `upstreamsById` для списка `model.upstreams`;
|
||||||
|
- **цикл по свободным апстримам** (наш внутренний учёт):
|
||||||
|
- `claim(...)` занимает слот первого свободного апстрима; если свободных нет —
|
||||||
|
сразу `503` (`all upstreams busy`) — решение по нашим счётчикам, не по апстриму;
|
||||||
|
- резолвим `providersById[upstream.provider]` (ошибка старта, если `id` нет);
|
||||||
|
- подменяем в теле `model` → `upstream.model`;
|
||||||
|
- **накладываем `patch` послойно** (исходное тело → `provider.patch` →
|
||||||
|
`upstream.patch` → `model.patch`), каждый через `merge(...)`; слои с `patch
|
||||||
|
== null` пропускаем. Итог — финальное тело запроса;
|
||||||
|
- шлём на `provider.url + /chat/completions` с `Authorization: Bearer provider.key`;
|
||||||
|
- при ответе `5xx`/`429` от провайдера — `release(upstream)` и переходим к
|
||||||
|
следующему свободному апстриму из списка (фейловер); исчерпали список —
|
||||||
|
отдаём `503` (`all upstreams failed`);
|
||||||
|
- при успехе/отмене клиента — `release(upstream)` в `finally` и выходим.
|
||||||
|
|
||||||
|
`patchBody(excluded, thinking)` и env-переменные `EXCLUDED_PROVIDERS` /
|
||||||
|
`THINKING_MODELS` / `UPSTREAM_URL` / `ROUTER_API_KEY` **убираются** — их
|
||||||
|
функциональность теперь в декларативном `patch` и блоке `upstreams`.
|
||||||
|
|
||||||
|
### 4. `handleModels` — отдаём свой каталог
|
||||||
|
|
||||||
|
Вместо проксирования `/v1/models` на апстрим — формируем ответ из
|
||||||
|
`config.models`, отдавая `id = model.name` для каждой записи. Суффикс `-no-think`
|
||||||
|
как отдельная модель теперь просто объявляется в конфиге (с нужным `patch`),
|
||||||
|
логика дублирования каталога удаляется.
|
||||||
|
|
||||||
|
### 5. Env, которые остаются
|
||||||
|
|
||||||
|
- `CONFIG_PATH` — путь к YAML; по умолчанию `config.yaml` в каталоге проекта
|
||||||
|
(CWD), можно задать абсолютный или относительный путь.
|
||||||
|
|
||||||
|
Порт/интерфейс биндинга больше не задаются через env — они в блоке `server`
|
||||||
|
файла конфига (`server.host`, `server.port`; дефолты `0.0.0.0` / `8100`).
|
||||||
|
|
||||||
|
### 6. Таймаут запроса к апстриму
|
||||||
|
|
||||||
|
Целое время на один запрос к нейронке (от отправки до получения всего ответа,
|
||||||
|
включая стриминг) — константа `UPSTREAM_REQUEST_TIMEOUT_MS` = **5 минут**
|
||||||
|
(`src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt`). Задается лимитом
|
||||||
|
`requestTimeout` CIO-движка Ktor — без этого работает дефолт движка **15 секунд**,
|
||||||
|
который обрывает длинные LLM-генерации. По достижении лимита запрос отменяется,
|
||||||
|
слот конкурентности освобождается. Значение выводится в лог при старте.
|
||||||
|
|
||||||
|
### 7. Логирование
|
||||||
|
|
||||||
|
Логируем через `logback` (зависимость `logback-classic` уже в проекте; `println`
|
||||||
|
заменяем на `Logger`). Ключевые события:
|
||||||
|
|
||||||
|
- **Старт**: сколько провайдеров/апстримов/моделей загружено; предупреждение,
|
||||||
|
если `upstreams[].id` или `models[].upstreams` ссылаются на несуществующий
|
||||||
|
`id` (проблема конфигурации).
|
||||||
|
- **Биндинг сервера**: `server: bind host=<host> port=<port>` — фактический
|
||||||
|
интерфейс и порт из блока `server` конфига.
|
||||||
|
- **Таймаут апстрима**: `upstream: request_timeout=<N>ms` — лимит времени на
|
||||||
|
запрос к нейронке (константа из раздела 6).
|
||||||
|
- **Принят запрос** (`model=<витрина>`, upstream=<id>, provider=<id>`): занят
|
||||||
|
слот `active[id]=N/limit`.
|
||||||
|
- **Отклонён запрос** — `503 all upstreams busy` для `model=<витрина>`:
|
||||||
|
состояние слотов всех апстримов модели (`id=N/limit`, ...).
|
||||||
|
- **Фейловер**: `upstream=<id> вернул <status>` → переход к следующему
|
||||||
|
свободному `upstream=<id>`.
|
||||||
|
- **Завершение** (успех/ошибка/отмена клиента): длительность, статус апстрима,
|
||||||
|
освобождён слот `active[id]=N/limit`.
|
||||||
|
- **Ошибка апстрима** (не `5xx`/`429`, а сетевая/таймаут): как сейчас — лог с
|
||||||
|
сообщением.
|
||||||
|
|
||||||
|
Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать
|
||||||
|
существующий парсинг логов (если он есть).
|
||||||
|
|
||||||
|
## Миграция `podman-compose.yaml`
|
||||||
|
|
||||||
|
Блок `environment` упрощается: вместо `UPSTREAM_URL` / `ROUTER_API_KEY` /
|
||||||
|
`EXCLUDED_PROVIDERS` / `THINKING_MODELS` монтируется файл конфига и задаётся
|
||||||
|
`CONFIG_PATH`, а секреты (ключи) — через env-подстановку `${...}` внутри файла.
|
||||||
@@ -2,8 +2,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.header
|
import io.ktor.client.request.headers
|
||||||
import io.ktor.utils.io.readAvailable
|
import io.ktor.utils.io.readAvailable
|
||||||
|
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
|
||||||
import io.ktor.client.statement.HttpResponse
|
import io.ktor.client.statement.HttpResponse
|
||||||
@@ -181,13 +182,24 @@ private suspend fun handleChat(
|
|||||||
}
|
}
|
||||||
|
|
||||||
val url = provider.url.trimEnd('/') + "/chat/completions"
|
val url = provider.url.trimEnd('/') + "/chat/completions"
|
||||||
|
val providerKey = resolveEnv(provider.key)
|
||||||
var failover = false
|
var failover = false
|
||||||
var responded = false
|
var responded = false
|
||||||
var upstreamStatus = 0
|
var upstreamStatus = 0
|
||||||
|
|
||||||
http.preparePost(url) {
|
http.preparePost(url) {
|
||||||
header("Authorization", "Bearer ${resolveEnv(provider.key)}")
|
headers {
|
||||||
header("Content-Type", "application/json")
|
headersToForward(call.request.headers).forEach { (name, values) ->
|
||||||
|
appendAll(name, values)
|
||||||
|
}
|
||||||
|
// Авторизация — всегда наша (ключ провайдера из конфига);
|
||||||
|
// клиентский Authorization не пересылается. Если у провайдера
|
||||||
|
// ключ не задан — Authorization не отправляем вовсе.
|
||||||
|
if (providerKey.isNotEmpty()) {
|
||||||
|
set("Authorization", "Bearer $providerKey")
|
||||||
|
}
|
||||||
|
set("Content-Type", "application/json")
|
||||||
|
}
|
||||||
setBody(forwarded.toString())
|
setBody(forwarded.toString())
|
||||||
}.execute { resp ->
|
}.execute { resp ->
|
||||||
upstreamStatus = resp.status.value
|
upstreamStatus = resp.status.value
|
||||||
@@ -262,6 +274,34 @@ private suspend fun handleChat(
|
|||||||
private fun errorJson(msg: String): String =
|
private fun errorJson(msg: String): String =
|
||||||
"""{"error":{"message":"$msg"}}"""
|
"""{"error":{"message":"$msg"}}"""
|
||||||
|
|
||||||
|
/** Hop-by-hop и прокси-специфичные заголовки, которые НЕ пересылаются в апстрим. */
|
||||||
|
private val SKIP_HEADER_NAMES = setOf(
|
||||||
|
"host",
|
||||||
|
"content-length",
|
||||||
|
"transfer-encoding",
|
||||||
|
"te",
|
||||||
|
"connection",
|
||||||
|
"proxy-connection",
|
||||||
|
"keep-alive",
|
||||||
|
"upgrade",
|
||||||
|
// Authorization управляется прокси явно (ключ провайдера), клиентский
|
||||||
|
// не пересылается
|
||||||
|
"authorization",
|
||||||
|
)
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Заголовки входящего запроса для пересылки в апстрим: все, кроме служебных
|
||||||
|
* ([SKIP_HEADER_NAMES]). Тело может быть изменено `patch`, а соединение до
|
||||||
|
* апстрима — другое, поэтому Content-Length/Transfer-Encoding/Host/Connection
|
||||||
|
* управляет сам прокси (клиентский Ktor выставит свои автоматически).
|
||||||
|
* `Authorization` тоже в списке пропусков — прокси управляет им явно
|
||||||
|
* (ключ провайдера из конфига; клиентский не пересылается).
|
||||||
|
*/
|
||||||
|
internal fun headersToForward(request: Headers): Map<String, List<String>> =
|
||||||
|
request.entries()
|
||||||
|
.filter { (name, _) -> name.lowercase() !in SKIP_HEADER_NAMES }
|
||||||
|
.associate { (name, values) -> name to values }
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Сборка тела запроса: подмена `model` на реальное имя апстрима + глубокий
|
* Сборка тела запроса: подмена `model` на реальное имя апстрима + глубокий
|
||||||
* послойный мерж `patch` в порядке provider → upstream → model.
|
* послойный мерж `patch` в порядке provider → upstream → model.
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import kotlin.test.Test
|
|||||||
import kotlin.test.assertEquals
|
import kotlin.test.assertEquals
|
||||||
import kotlin.test.assertFalse
|
import kotlin.test.assertFalse
|
||||||
import kotlin.test.assertTrue
|
import kotlin.test.assertTrue
|
||||||
|
import io.ktor.http.headersOf
|
||||||
import kotlinx.serialization.json.Json
|
import kotlinx.serialization.json.Json
|
||||||
import kotlinx.serialization.json.JsonObject
|
import kotlinx.serialization.json.JsonObject
|
||||||
import kotlinx.serialization.json.jsonArray
|
import kotlinx.serialization.json.jsonArray
|
||||||
@@ -275,6 +276,42 @@ class ConfigLogicTest {
|
|||||||
assertEquals(null, pickFreeUpstream(pool, active, emptySet())?.id)
|
assertEquals(null, pickFreeUpstream(pool, active, emptySet())?.id)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun headersToForwardKeepsUnknownAndDropsServiceHeaders() {
|
||||||
|
val req = headersOf(
|
||||||
|
"X-Opencode-Session" to listOf("abc-123"),
|
||||||
|
"X-Custom" to listOf("a", "b"),
|
||||||
|
"Host" to listOf("proxy:8100"),
|
||||||
|
"Content-Length" to listOf("42"),
|
||||||
|
"Transfer-Encoding" to listOf("chunked"),
|
||||||
|
"Connection" to listOf("keep-alive"),
|
||||||
|
"TE" to listOf("trailers"),
|
||||||
|
"Proxy-Connection" to listOf("keep-alive"),
|
||||||
|
"Upgrade" to listOf("h2c"),
|
||||||
|
"Authorization" to listOf("Bearer client-secret"),
|
||||||
|
"Accept" to listOf("*/*"),
|
||||||
|
)
|
||||||
|
val out = headersToForward(req)
|
||||||
|
assertEquals(listOf("abc-123"), out["X-Opencode-Session"])
|
||||||
|
assertEquals(listOf("a", "b"), out["X-Custom"])
|
||||||
|
assertEquals(listOf("*/*"), out["Accept"])
|
||||||
|
assertEquals(3, out.size)
|
||||||
|
assertEquals(null, out["Authorization"])
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun headersToForwardIsCaseInsensitiveOnSkipSet() {
|
||||||
|
val out = headersToForward(
|
||||||
|
headersOf(
|
||||||
|
"HOST" to listOf("x"),
|
||||||
|
"Content-length" to listOf("1"),
|
||||||
|
"x-opencode-session" to listOf("s"),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
assertEquals(1, out.size)
|
||||||
|
assertEquals(listOf("s"), out["x-opencode-session"])
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun rebuildFromChunksPreservesAllUpstreamFields() {
|
fun rebuildFromChunksPreservesAllUpstreamFields() {
|
||||||
val sse = """
|
val sse = """
|
||||||
|
|||||||
Reference in New Issue
Block a user