Compare commits
4 Commits
6
...
dev-20260911-0527
| Author | SHA1 | Date | |
|---|---|---|---|
| dba52d22fd | |||
| efdce75ee9 | |||
| e251ca5326 | |||
| 0fb99060f9 |
@@ -5,3 +5,9 @@ build/
|
||||
# локальный конфиг с секретами (ключ API) — не коммитим
|
||||
config.yaml
|
||||
|
||||
# IDE / tooling
|
||||
.idea/
|
||||
.kotlin/
|
||||
.cortexkit/
|
||||
.veai/
|
||||
|
||||
|
||||
@@ -0,0 +1,426 @@
|
||||
# Конфигурация 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 # опционально; лимит по умолчанию для апстримов
|
||||
session_header: x-opencode-session # опционально; прокси считает сессию из истории
|
||||
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`):
|
||||
если ни одного нет — запрос проксируется как есть (исходное тело клиента).
|
||||
|
||||
### Сессия по истории (`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`-сообщением неразличимы (склеятся).
|
||||
|
||||
### Пример сборки тела (многослойный `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, // ко всем запросам провайдера
|
||||
val session_header: String? = 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.call.body
|
||||
import io.ktor.client.request.header
|
||||
import io.ktor.client.request.headers
|
||||
import io.ktor.utils.io.readAvailable
|
||||
import io.ktor.http.Headers
|
||||
import io.ktor.client.request.preparePost
|
||||
import io.ktor.client.request.setBody
|
||||
import io.ktor.client.statement.HttpResponse
|
||||
@@ -13,7 +14,9 @@ import io.ktor.server.application.Application
|
||||
import io.ktor.server.application.ApplicationCall
|
||||
import io.ktor.server.application.call
|
||||
import io.ktor.server.application.install
|
||||
import io.ktor.server.request.httpMethod
|
||||
import io.ktor.server.request.receiveText
|
||||
import io.ktor.server.request.uri
|
||||
import io.ktor.server.response.respondBytesWriter
|
||||
import io.ktor.server.response.respondText
|
||||
import io.ktor.server.routing.get
|
||||
@@ -94,8 +97,9 @@ fun main() {
|
||||
}
|
||||
|
||||
val http = createHttpClient()
|
||||
val sessions = SessionRegistry()
|
||||
startServer(config.server.host, config.server.port) {
|
||||
proxyModule(config, providersById, upstreamsById, active, http)
|
||||
proxyModule(config, providersById, upstreamsById, active, sessions, http)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -106,11 +110,12 @@ fun Application.proxyModule(
|
||||
providersById: Map<String, ProviderConf>,
|
||||
upstreamsById: Map<String, UpstreamConf>,
|
||||
active: Map<String, UpstreamCounter>,
|
||||
sessions: SessionRegistry,
|
||||
http: HttpClient,
|
||||
) {
|
||||
routing {
|
||||
post("/v1/chat/completions") {
|
||||
handleChat(call, config, providersById, upstreamsById, active, http)
|
||||
handleChat(call, config, providersById, upstreamsById, active, sessions, http)
|
||||
}
|
||||
get("/v1/models") {
|
||||
handleModels(call, config)
|
||||
@@ -124,8 +129,10 @@ private suspend fun handleChat(
|
||||
providersById: Map<String, ProviderConf>,
|
||||
upstreamsById: Map<String, UpstreamConf>,
|
||||
active: Map<String, UpstreamCounter>,
|
||||
sessions: SessionRegistry,
|
||||
http: HttpClient,
|
||||
) {
|
||||
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
|
||||
val raw = call.receiveText()
|
||||
if (raw.isBlank()) {
|
||||
call.respondText(errorJson("empty body"), ContentType.Application.Json, HttpStatusCode.BadRequest)
|
||||
@@ -181,13 +188,41 @@ private suspend fun handleChat(
|
||||
}
|
||||
|
||||
val url = provider.url.trimEnd('/') + "/chat/completions"
|
||||
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 responded = false
|
||||
var upstreamStatus = 0
|
||||
|
||||
http.preparePost(url) {
|
||||
header("Authorization", "Bearer ${resolveEnv(provider.key)}")
|
||||
header("Content-Type", "application/json")
|
||||
headers {
|
||||
forwardedHeaders.forEach { (name, values) ->
|
||||
appendAll(name, values)
|
||||
}
|
||||
// Авторизация — всегда наша (ключ провайдера из конфига);
|
||||
// клиентский Authorization не пересылается. Если у провайдера
|
||||
// ключ не задан — Authorization не отправляем вовсе.
|
||||
if (providerKey.isNotEmpty()) {
|
||||
set("Authorization", "Bearer $providerKey")
|
||||
}
|
||||
set("Content-Type", "application/json")
|
||||
if (sessionHeader != null && sessionId != null) {
|
||||
set(sessionHeader, sessionId)
|
||||
}
|
||||
}
|
||||
setBody(forwarded.toString())
|
||||
}.execute { resp ->
|
||||
upstreamStatus = resp.status.value
|
||||
@@ -262,6 +297,60 @@ private suspend fun handleChat(
|
||||
private fun errorJson(msg: String): String =
|
||||
"""{"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",
|
||||
// Content-Type всегда наш (application/json: тело мержится как JSON),
|
||||
// клиентский не пересылаем, чтобы не ушло двух заголовков
|
||||
"content-type",
|
||||
)
|
||||
|
||||
/**
|
||||
* Заголовки входящего запроса для пересылки в апстрим: все, кроме служебных
|
||||
* ([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 }
|
||||
|
||||
/** Заголовки, значения которых маскируются в логах (секреты клиента). */
|
||||
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` на реальное имя апстрима + глубокий
|
||||
* послойный мерж `patch` в порядке provider → upstream → model.
|
||||
@@ -382,6 +471,7 @@ internal fun pickFreeUpstream(
|
||||
pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) }
|
||||
|
||||
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 data = config.models.map { m ->
|
||||
JsonObject(
|
||||
@@ -497,6 +587,7 @@ data class ProviderConf(
|
||||
val key: String = "",
|
||||
val max_concurrency: Int? = null,
|
||||
val patch: JsonObject? = null,
|
||||
val session_header: String? = null,
|
||||
)
|
||||
|
||||
data class UpstreamConf(
|
||||
@@ -558,6 +649,7 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
key = m.strOrNull("key") ?: "",
|
||||
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
|
||||
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
|
||||
session_header = m.strOrNull("session_header"),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
import io.ktor.http.headersOf
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.json.JsonObject
|
||||
import kotlinx.serialization.json.jsonArray
|
||||
@@ -275,6 +276,70 @@ class ConfigLogicTest {
|
||||
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"),
|
||||
"Content-Type" to listOf("application/x-www-form-urlencoded"),
|
||||
"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"])
|
||||
assertEquals(null, out["Content-Type"])
|
||||
}
|
||||
|
||||
@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
|
||||
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
|
||||
fun rebuildFromChunksPreservesAllUpstreamFields() {
|
||||
val sse = """
|
||||
@@ -317,4 +382,116 @@ class ConfigLogicTest {
|
||||
val err = out["error"]?.jsonObject ?: error("error block missing")
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user