feat: Prometheus-метрики GET /metrics + ленивый ISO-8601-парсер для backoff
Build LLM Proxy / Build and push (release) Successful in 40s
Build LLM Proxy / Build and push (release) Successful in 40s
- GET /metrics (Prometheus text 0.0.4): llm_proxy_requests_total{model,upstream,provider,result}
(ok/4xx/429/402/5xx/net_err/cancelled), llm_proxy_upstream_inflight,
llm_proxy_upstream_fail_streak, llm_proxy_upstream_cooling_seconds
- учёт исходов в handleChat (каждая фейловер-попытка — отдельно)
- parseIsoDuration: PnM≈n×30d, PnY≈n×365d (kotlin.time Duration.parse не принимает Y/M)
- remainingSeconds округляет остаток вверх (открытое окно ≥1s)
- доки (CONFIG.md/README/TESTING) + тесты (MetricsTest, BackoffTest)
This commit is contained in:
@@ -297,10 +297,11 @@ providers:
|
||||
(`all upstreams cooling, retry in Ns`) с заголовком `Retry-After: N`
|
||||
(секунд до выхода первого апстрима из отката).
|
||||
|
||||
Значение — ISO-8601-длительность (формат, в котором сериализуется
|
||||
`kotlin.time.Duration`): `PT30M` (30 минут), `PT1H15M`, `P1D` (сутки),
|
||||
`P1M` (месяц). Некорректное значение — предупреждение в лог и поле
|
||||
просто игнорируется.
|
||||
Значение — ISO-8601-длительность: `PT30M` (30 минут), `PT1H15M`, `P1D` (сутки),
|
||||
`P1M` (месяц), `P1Y` (год). Годы/месяцы укорачиваются приближённо
|
||||
(1 год ≈ 365d, 1 месяц ≈ 30d) — kotlin.time `Duration.parse` не принимает
|
||||
Y/M (у них нет фиксированной длины), конфиг-парсер прокси расширяет формат.
|
||||
Некорректное значение — предупреждение в лог и поле просто игнорируется.
|
||||
|
||||
```yaml
|
||||
providers:
|
||||
@@ -318,6 +319,29 @@ upstreams:
|
||||
При старте выводится, что настроено:
|
||||
`[llm-proxy] backoff: upstreams=routerai-gpt4o=PT15M providers=routerai=P1M`.
|
||||
|
||||
### Prometheus-метрики (`/metrics`)
|
||||
|
||||
Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics`
|
||||
(`text/plain; version=0.0.4`), без авторизации (внутренний контур).
|
||||
Достаточно включить скрейп в Prometheus/VictoriaMetrics — и в Grafana можно
|
||||
вести дашборды использования по провайдерам/моделям и алерты на деградацию.
|
||||
|
||||
| Метрика | Тип | Смысл |
|
||||
|---|---|---|
|
||||
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `cancelled` — клиент отвалился посреди стрима |
|
||||
| `llm_proxy_upstream_inflight{upstream, provider}` | gauge | занятые слоты апстрима прямо сейчас (конкурентность) |
|
||||
| `llm_proxy_upstream_fail_streak{upstream, provider}` | gauge | счётчик сбоев подряд (backoff): сколько раз подряд упал |
|
||||
| `llm_proxy_upstream_cooling_seconds{upstream, provider}` | gauge | сколько секунд апстрим ещё в backoff-откате (0 = жив) |
|
||||
|
||||
Метки: `model` — витринное имя модели, `upstream`/`provider` — внутренние id из конфига.
|
||||
Метрики live в памяти: при рестарте прокси сбрасываются (серить их будет Prometheus).
|
||||
|
||||
Примеры для Grafana:
|
||||
- оборот по виртуальным моделям: `sum(rate(llm_proxy_requests_total[5m])) by (model)`;
|
||||
- оборот по провайдерам: `sum(rate(llm_proxy_requests_total[5m])) by (provider)`;
|
||||
- «провайдер умер»: `llm_proxy_upstream_cooling_seconds > 0` дольше N минут — алерт;
|
||||
- доля ошибок провайдера: `sum(rate(llm_proxy_requests_total{result=~"4xx|429|402|5xx|net_err"}[10m])) by (provider) / sum(rate(llm_proxy_requests_total[10m])) by (provider)`.
|
||||
|
||||
### Пример сборки тела (многослойный `patch`)
|
||||
|
||||
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
|
||||
|
||||
@@ -30,6 +30,11 @@ non-stream.
|
||||
1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате
|
||||
— `503` с `Retry-After`.
|
||||
|
||||
Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по
|
||||
`model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) —
|
||||
подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список
|
||||
метрик — в CONFIG.md, раздел «Prometheus-метрики».
|
||||
|
||||
| Переменная | Default | Описание |
|
||||
|---|---|---|
|
||||
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
|
||||
|
||||
+2
-1
@@ -20,7 +20,8 @@
|
||||
| `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; |
|
||||
| `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. |
|
||||
| `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. |
|
||||
| `BackoffTest` | Экспоненциальный backoff: удвоение интервала отката до потолка (cap), сброс счётчика при успехе, окно охлаждения (`isCoolingAt`/`remainingAt`), ISO-8601-разбор cap (`PT30M`, `P1D`, `PT1M30S`), скоупы: общий провайдерский счётчик на все его модели, приоритет `upstreams[].backoff` над провайдерским, независимость скоупов, backoff не настроен → без отката. |
|
||||
| `BackoffTest` | Экспоненциальный backoff: удвоение интервала отката до потолка (cap), сброс счётчика при успехе, окно охлаждения (`isCoolingAt`/`remainingAt`), ISO-8601-разбор cap (`PT30M`, `P1D`, `PT1M30S`) и конфиг-парсер `parseIsoDuration` (ленивые Y/M: `P1M`≈30d, `P1Y`≈365d; мусор — ошибка), скоупы: общий провайдерский счётчик на все его модели, приоритет `upstreams[].backoff` над провайдерским, независимость скоупов, backoff не настроен → без отката. |
|
||||
| `MetricsTest` | `/metrics`: counters по исходам запросов (`ok/4xx/429/402/5xx/net_err/cancelled`), гейджи `inflight`/`fail_streak`/`cooling_seconds` (общий провайдерский скоуп виден в метриках), экранирование меток, маппинг HTTP-статуса → `result`. |
|
||||
|
||||
## Проверка качества тестов (мутационная приёмка)
|
||||
|
||||
|
||||
@@ -52,6 +52,9 @@ class BackoffGuard(val cap: Duration, val base: Duration = 1.seconds) {
|
||||
lock.unlock()
|
||||
}
|
||||
}
|
||||
|
||||
/** Текущая серия сбоев подряд (0 — после успеха или без сбоев). */
|
||||
fun failStreak(): Int = failures
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -92,10 +95,14 @@ class BackoffRegistry(
|
||||
return g.isCoolingAt(Clock.System.now())
|
||||
}
|
||||
|
||||
/** Секунд до выхода апстрима из отдыха (0 — не в откате). */
|
||||
/** Секунд до выхода апстрима из отдыха (0 — не в откате); дробный остаток округляется вверх. */
|
||||
suspend fun remainingSeconds(up: UpstreamConf): Long {
|
||||
val g = locked { guardForLocked(up) } ?: return 0L
|
||||
return g.remainingAt(Clock.System.now()).inWholeSeconds
|
||||
val remaining = g.remainingAt(Clock.System.now())
|
||||
if (remaining <= Duration.ZERO) return 0L
|
||||
val s = remaining.inWholeSeconds
|
||||
val isExact = remaining.inWholeNanoseconds == s * 1_000_000_000L
|
||||
return if (isExact) s else s + 1
|
||||
}
|
||||
|
||||
/** Учёт ошибки; возвращает интервал отдыха, или null, если бэкофф для апстрима выключен. */
|
||||
@@ -108,4 +115,10 @@ class BackoffRegistry(
|
||||
suspend fun recordSuccess(up: UpstreamConf) {
|
||||
locked { guardForLocked(up) }?.recordSuccess()
|
||||
}
|
||||
|
||||
/** Текущая серия сбоев апстрима (0 — нет сбоев или backoff не настроен). */
|
||||
suspend fun failStreakOf(up: UpstreamConf): Int {
|
||||
val g = locked { guardForLocked(up) } ?: return 0
|
||||
return g.failStreak()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -47,6 +47,9 @@ import net.mamoe.yamlkt.YamlLiteral
|
||||
import net.mamoe.yamlkt.YamlMap
|
||||
import kotlin.concurrent.Volatile
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.days
|
||||
import kotlin.time.Duration.Companion.hours
|
||||
import kotlin.time.Duration.Companion.minutes
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
import kotlin.time.TimeSource
|
||||
|
||||
@@ -88,6 +91,7 @@ fun main() {
|
||||
up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider]))
|
||||
}
|
||||
val backoff = BackoffRegistry(providersById, upstreamsById)
|
||||
val metrics = MetricsRegistry()
|
||||
|
||||
log.info {
|
||||
"[llm-proxy] загружено: providers=${config.providers.size}, " +
|
||||
@@ -109,7 +113,7 @@ fun main() {
|
||||
val http = createHttpClient()
|
||||
val sessions = SessionRegistry()
|
||||
startServer(config.server.host, config.server.port) {
|
||||
proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff)
|
||||
proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -123,10 +127,18 @@ fun Application.proxyModule(
|
||||
sessions: SessionRegistry,
|
||||
http: HttpClient,
|
||||
backoff: BackoffRegistry,
|
||||
metrics: MetricsRegistry,
|
||||
) {
|
||||
routing {
|
||||
get("/metrics") {
|
||||
call.respondText(
|
||||
metrics.render(config, active, backoff),
|
||||
ContentType.parse("text/plain; version=0.0.4"),
|
||||
HttpStatusCode.OK,
|
||||
)
|
||||
}
|
||||
post("/v1/chat/completions") {
|
||||
handleChat(call, config, providersById, upstreamsById, active, sessions, http, backoff)
|
||||
handleChat(call, config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
|
||||
}
|
||||
get("/v1/models") {
|
||||
handleModels(call, config)
|
||||
@@ -143,6 +155,7 @@ private suspend fun handleChat(
|
||||
sessions: SessionRegistry,
|
||||
http: HttpClient,
|
||||
backoff: BackoffRegistry,
|
||||
metrics: MetricsRegistry,
|
||||
) {
|
||||
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
|
||||
val raw = call.receiveText()
|
||||
@@ -241,6 +254,7 @@ private suspend fun handleChat(
|
||||
upstreamStatus = resp.status.value
|
||||
if (upstreamStatus >= 500 || upstreamStatus == 429 || upstreamStatus == 402) {
|
||||
val rest = backoff.recordFailure(up)
|
||||
metrics.record(modelName, up, statusClass(upstreamStatus))
|
||||
log.warn {
|
||||
"[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" +
|
||||
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
|
||||
@@ -252,7 +266,8 @@ private suspend fun handleChat(
|
||||
if (upstreamStatus in 400..499) {
|
||||
// 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть.
|
||||
val errorBody = runCatching { resp.body<String>() }.getOrDefault("")
|
||||
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" }
|
||||
log.warn { "[llm-proxy] chat model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" }
|
||||
metrics.record(modelName, up, statusClass(upstreamStatus))
|
||||
val ct = resp.headers["Content-Type"] ?: "application/json"
|
||||
call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus))
|
||||
responded = true
|
||||
@@ -273,6 +288,7 @@ private suspend fun handleChat(
|
||||
streamSseWithThinkTags(ch, thinkMode)
|
||||
}
|
||||
}
|
||||
metrics.record(modelName, up, "ok")
|
||||
log.info {
|
||||
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
|
||||
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=true"
|
||||
@@ -296,6 +312,7 @@ private suspend fun handleChat(
|
||||
}
|
||||
}.getOrDefault(ContentType.parse(ct))
|
||||
call.respondText(out, outCt, HttpStatusCode.fromValue(upstreamStatus))
|
||||
metrics.record(modelName, up, "ok")
|
||||
log.info {
|
||||
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
|
||||
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=false"
|
||||
@@ -306,10 +323,12 @@ private suspend fun handleChat(
|
||||
if (failover) continue
|
||||
if (responded) return
|
||||
} catch (e: CancellationException) {
|
||||
metrics.record(modelName, up, "cancelled")
|
||||
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
val rest = backoff.recordFailure(up)
|
||||
metrics.record(modelName, up, "net_err")
|
||||
log.error {
|
||||
"[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" +
|
||||
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
|
||||
@@ -1052,14 +1071,43 @@ internal fun Map<String, YamlElement>.str(key: String): String =
|
||||
internal fun Map<String, YamlElement>.strOrNull(key: String): String? =
|
||||
(this[key] as? YamlLiteral)?.content
|
||||
|
||||
/** ISO-8601-длительность (PT30M, P1M, PT1H15M…); некорректное значение — warn + null. */
|
||||
/**
|
||||
* ISO-8601-длительность (PT30M, P1D, PT1H15M…) с расширением для годов/месяцев,
|
||||
* которые kotlin.time не принимает (у них нет фиксированной длины):
|
||||
* `PnY` → n×365d, `PnM` → n×30d (приближение, см. CONFIG.md).
|
||||
* Некорректное значение — warn + null.
|
||||
*/
|
||||
internal fun Map<String, YamlElement>.durationOrNull(key: String): Duration? =
|
||||
strOrNull(key)?.let { raw ->
|
||||
runCatching { Duration.parse(raw) }.getOrElse {
|
||||
runCatching { parseIsoDuration(raw) }.getOrElse {
|
||||
log.warn { "[llm-proxy] config: поле '$key' — некорректная ISO-8601-длительность '$raw', проигнорировано" }
|
||||
null
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* ISO-8601-длительность для конфига: как [Duration.parse], плюс годы/месяцы
|
||||
* (`P1M`, `P2Y`, `P1Y3M`) — приближённо 1 год = 365d, 1 месяц = 30d.
|
||||
* (kotlin.time [Duration.parse] принимает только D/W/H/M(ин)/S.)
|
||||
*/
|
||||
internal fun parseIsoDuration(raw: String): Duration {
|
||||
val m = Regex(
|
||||
"^P" +
|
||||
"(?:(\\d+)Y)?" +
|
||||
"(?:(\\d+)M)?" +
|
||||
"(?:(\\d+)W)?" +
|
||||
"(?:(\\d+)D)?" +
|
||||
"(?:T(?:(\\d+)H)?(?:(\\d+)M)?(?:(\\d+)S)?)?",
|
||||
RegexOption.IGNORE_CASE,
|
||||
).matchEntire(raw.trim()) ?: error("не ISO-8601-длительность: $raw")
|
||||
if ((1..7).all { m.groupValues[it].isEmpty() }) error("не ISO-8601-длительность: $raw")
|
||||
fun g(i: Int): Long = m.groupValues[i].toLongOrNull() ?: 0L
|
||||
val days = g(1) * 365 + g(2) * 30 + g(3) * 7 + g(4)
|
||||
val hours = g(5)
|
||||
val minutes = g(6)
|
||||
val seconds = g(7)
|
||||
return days.days + hours.hours + minutes.minutes + seconds.seconds
|
||||
}
|
||||
|
||||
internal fun Map<String, YamlElement>.yamlMapOrNull(key: String): YamlMap? =
|
||||
this[key] as? YamlMap
|
||||
|
||||
@@ -0,0 +1,100 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
|
||||
/**
|
||||
* Pull-метрики для Prometheus (формат экспозиции `text/plain; version=0.0.4`),
|
||||
* отдаются на `GET /metrics`.
|
||||
*
|
||||
* Метрики (всё в памяти, сброс при рестарте):
|
||||
* - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик
|
||||
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | cancelled`;
|
||||
* - `llm_proxy_upstream_inflight{upstream,provider}` — сейчас в работе (слоты);
|
||||
* - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев
|
||||
* (счётчик backoff);
|
||||
* - `llm_proxy_upstream_cooling_seconds{upstream,provider}` — сколько секунд
|
||||
* апстрим ещё в backoff-откате (0 = жив).
|
||||
*/
|
||||
class MetricsRegistry {
|
||||
private val lock = Mutex()
|
||||
private val requests = HashMap<RequestKey, Long>()
|
||||
|
||||
private data class RequestKey(
|
||||
val model: String,
|
||||
val upstream: String,
|
||||
val provider: String,
|
||||
val result: String,
|
||||
)
|
||||
|
||||
/** Учёт chat-запроса по исходу (значения `result` — в описании класса). */
|
||||
suspend fun record(model: String, up: UpstreamConf, result: String) {
|
||||
lock.withLock {
|
||||
val key = RequestKey(model, up.id, up.provider, result)
|
||||
requests[key] = (requests[key] ?: 0L) + 1
|
||||
}
|
||||
}
|
||||
|
||||
/** Текстовая выгрузка в формате Prometheus на текущий момент. */
|
||||
suspend fun render(
|
||||
config: Config,
|
||||
active: Map<String, UpstreamCounter>,
|
||||
backoff: BackoffRegistry,
|
||||
): String {
|
||||
val sb = StringBuilder()
|
||||
sb.append("# HELP llm_proxy_requests_total Chat-запросы llm-proxy по исходу\n")
|
||||
sb.append("# TYPE llm_proxy_requests_total counter\n")
|
||||
lock.withLock {
|
||||
requests.entries
|
||||
.sortedWith(compareBy({ it.key.model }, { it.key.upstream }, { it.key.provider }, { it.key.result }))
|
||||
.forEach { (k, n) ->
|
||||
sb.append(
|
||||
"llm_proxy_requests_total{model=\"" + esc(k.model) +
|
||||
"\",upstream=\"" + esc(k.upstream) +
|
||||
"\",provider=\"" + esc(k.provider) +
|
||||
"\",result=\"" + esc(k.result) + "\"} $n\n",
|
||||
)
|
||||
}
|
||||
}
|
||||
sb.append("# HELP llm_proxy_upstream_inflight Занято слотов апстрима сейчас\n")
|
||||
sb.append("# TYPE llm_proxy_upstream_inflight gauge\n")
|
||||
for (up in config.upstreams) {
|
||||
val n = active[up.id]?.current ?: 0
|
||||
sb.append(
|
||||
"llm_proxy_upstream_inflight{upstream=\"" + esc(up.id) +
|
||||
"\",provider=\"" + esc(up.provider) + "\"} $n\n",
|
||||
)
|
||||
}
|
||||
sb.append("# HELP llm_proxy_upstream_fail_streak Серия сбоев апстрима подряд (счётчик backoff)\n")
|
||||
sb.append("# TYPE llm_proxy_upstream_fail_streak gauge\n")
|
||||
for (up in config.upstreams) {
|
||||
sb.append(
|
||||
"llm_proxy_upstream_fail_streak{upstream=\"" + esc(up.id) +
|
||||
"\",provider=\"" + esc(up.provider) + "\"} " +
|
||||
backoff.failStreakOf(up) + "\n",
|
||||
)
|
||||
}
|
||||
sb.append("# HELP llm_proxy_upstream_cooling_seconds Секунд до выхода апстрима из backoff-отката (0 = не в откате)\n")
|
||||
sb.append("# TYPE llm_proxy_upstream_cooling_seconds gauge\n")
|
||||
for (up in config.upstreams) {
|
||||
sb.append(
|
||||
"llm_proxy_upstream_cooling_seconds{upstream=\"" + esc(up.id) +
|
||||
"\",provider=\"" + esc(up.provider) + "\"} " +
|
||||
backoff.remainingSeconds(up) + "\n",
|
||||
)
|
||||
}
|
||||
return sb.toString()
|
||||
}
|
||||
|
||||
private fun esc(s: String): String =
|
||||
s.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n")
|
||||
}
|
||||
|
||||
/** HTTP-статус апстрима → класс исхода для метрик. */
|
||||
internal fun statusClass(status: Int): String = when {
|
||||
status in 200..399 -> "ok"
|
||||
status == 429 -> "429"
|
||||
status == 402 -> "402"
|
||||
status in 400..499 -> "4xx"
|
||||
else -> "5xx"
|
||||
}
|
||||
@@ -4,11 +4,13 @@ import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.datetime.Instant
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFails
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.days
|
||||
import kotlin.time.Duration.Companion.hours
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.minutes
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
@@ -53,6 +55,26 @@ class BackoffTest {
|
||||
assertEquals(90.seconds, Duration.parse("PT1M30S"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun configDurationParsingWithYearsAndMonths() {
|
||||
// kotlin.time не принимает Y/M — конфиг-парсер расширяет: 1Y≈365d, 1M≈30d
|
||||
assertEquals(30.days, parseIsoDuration("P1M"))
|
||||
assertEquals(365.days, parseIsoDuration("P1Y"))
|
||||
assertEquals(455.days, parseIsoDuration("P1Y3M"))
|
||||
assertEquals(7.days, parseIsoDuration("P1W"))
|
||||
assertEquals(2.days + 3.hours + 15.minutes, parseIsoDuration("P2DT3H15M"))
|
||||
// базовый ISO-8601 без Y/M — как kotlin.time
|
||||
assertEquals(30.minutes, parseIsoDuration("PT30M"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun configDurationParsingRejectsGarbage() {
|
||||
assertFails { parseIsoDuration("PT30") }
|
||||
assertFails { parseIsoDuration("1d") }
|
||||
assertFails { parseIsoDuration("P") }
|
||||
assertFails { parseIsoDuration("") }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun providerScopeSharedBetweenItsModels() = runTest {
|
||||
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
class MetricsTest {
|
||||
|
||||
private val p = ProviderConf(id = "p", url = "https://p", backoff = 30.seconds)
|
||||
private val a = UpstreamConf(id = "a", provider = "p", model = "m1")
|
||||
private val b = UpstreamConf(id = "b", provider = "p", model = "m2")
|
||||
private val cfg = Config(
|
||||
server = ServerConf(),
|
||||
providers = listOf(p),
|
||||
upstreams = listOf(a, b),
|
||||
models = emptyList(),
|
||||
)
|
||||
|
||||
@Test
|
||||
fun requestCounters() = runTest {
|
||||
val m = MetricsRegistry()
|
||||
m.record("my-gpt", a, "ok")
|
||||
m.record("my-gpt", a, "ok")
|
||||
m.record("my-gpt", a, "429")
|
||||
m.record("my-gpt", b, "5xx")
|
||||
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4), "b" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
|
||||
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"a\",provider=\"p\",result=\"ok\"} 2\n"), "ok=2: $text")
|
||||
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"a\",provider=\"p\",result=\"429\"} 1\n"), "429=1: $text")
|
||||
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"b\",provider=\"p\",result=\"5xx\"} 1\n"), "5xx=1: $text")
|
||||
assertTrue(text.contains("# TYPE llm_proxy_requests_total counter"), "TYPE: $text")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun inflightGauge() = runTest {
|
||||
val m = MetricsRegistry()
|
||||
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4, 3), "b" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
|
||||
assertTrue(text.contains("llm_proxy_upstream_inflight{upstream=\"a\",provider=\"p\"} 3\n"), text)
|
||||
assertTrue(text.contains("llm_proxy_upstream_inflight{upstream=\"b\",provider=\"p\"} 0\n"), text)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun backoffGauges() = runTest {
|
||||
val backoff = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
|
||||
assertEquals(1.seconds, backoff.recordFailure(a))
|
||||
val text = MetricsRegistry().render(cfg, mapOf("a" to UpstreamCounter(4), "b" to UpstreamCounter(4)), backoff)
|
||||
// общий скоуп провайдера: сбой у a виден и у b
|
||||
assertTrue(text.contains("llm_proxy_upstream_fail_streak{upstream=\"a\",provider=\"p\"} 1\n"), text)
|
||||
assertTrue(text.contains("llm_proxy_upstream_fail_streak{upstream=\"b\",provider=\"p\"} 1\n"), text)
|
||||
assertTrue(Regex("llm_proxy_upstream_cooling_seconds\\{upstream=\"a\",provider=\"p\"\\} (0|1)").containsMatchIn(text), text)
|
||||
assertTrue(Regex("llm_proxy_upstream_cooling_seconds\\{upstream=\"b\",provider=\"p\"\\} (0|1)").containsMatchIn(text), text)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun labelEscaping() = runTest {
|
||||
val m = MetricsRegistry()
|
||||
m.record("we\"ird\nmodel", a, "ok")
|
||||
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
|
||||
assertTrue(text.contains("model=\"we\\\"ird\\nmodel\""), text)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun statusClasses() {
|
||||
assertEquals("ok", statusClass(200))
|
||||
assertEquals("ok", statusClass(302))
|
||||
assertEquals("4xx", statusClass(400))
|
||||
assertEquals("4xx", statusClass(499))
|
||||
assertEquals("429", statusClass(429))
|
||||
assertEquals("402", statusClass(402))
|
||||
assertEquals("5xx", statusClass(500))
|
||||
assertEquals("5xx", statusClass(599))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user