2 Commits

9 changed files with 917 additions and 11 deletions
+141 -4
View File
@@ -85,6 +85,8 @@ providers:
allow_fallbacks: false
backoff: P1M # опционально; потолок экспоненциального backoff
# на весь провайдер (ISO-8601: P1M, PT15M, P1D…)
probe_interval: PT1S # опционально; фоновый пинг апстримов в откате
# (как часто проверять, жив ли провайдер)
- id: local-llama
url: "http://10.0.0.5:8080/v1"
@@ -103,6 +105,8 @@ upstreams:
ignore: [deepseek]
backoff: PT15M # опционально; собственный backoff этой модели
# (приоритет над backoff провайдера)
probe_interval: PT0S # опционально; пинг отключён у этой модели
# (приоритет над probe_interval провайдера)
- id: local-qwen
provider: local-llama
@@ -322,6 +326,131 @@ upstreams:
При старте выводится, что настроено:
`[llm-proxy] backoff: upstreams=routerai-gpt4o=PT15M providers=routerai=P1M`.
### Таймаут первого байта стрима (`first_byte_timeout`)
Для `stream=true` запросов прокси следит, чтобы апстрим начал отдавать
**что-нибудь** (первый байт тела) в течение заданного окна. Если апстрим
молчит дольше — это трактуется как его сбой: слот конкурентности
освобождается, ошибка идёт в backoff-откат, и роутер **фейловерит** на
следующего свободного апстрима модели (при его наличии). Это закрывает
кейс, когда апстрим принял запрос, но «завис» и никогда не начинает
стримить.
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].first_byte_timeout` | ISO-8601-длительность / отсутствует | Окно ожидания первого байта для всех моделей провайдера |
| `upstreams[].first_byte_timeout` | ISO-8601-длительность / отсутствует | Окно для конкретной модели (приоритет над провайдерским) |
**Дефолт** (если не задан нигде): `PT30S`.
**Приоритет — модели.** Если `upstreams[].first_byte_timeout` задан, он
перекрывает провайдерский; если нет — берётся провайдерский; если нет и там —
дефолт `PT30S`.
**Отключение:** `PT0S` на провайдере или апстриме отключает watchdog
для этого апстрима — тогда первый байт ждём без лимита (в рамках общего
`upstream: request_timeout`, см. ниже).
**Применяется только к `stream=true`.** Для `stream=false` (не-стриминговый
ответ) действует только общий лимит `UPSTREAM_REQUEST_TIMEOUT_MS` (`5m`) —
строки SSE там склеиваются в один ответ, и таймаут первого байта к нему
неприменим.
Поведение:
- watchdog висит **на первом чтении** из тела стрим-ответа; как только
пришёл хоть какой-то байт — лимит отключается, дальше стрим читается
свободно (вплоть до общего request_timeout);
- `data: [DONE]` приходит в теле — это уже «первый байт», таймаут не
сработает;
- при срабатывании — лог:
`[llm-proxy] chat model=<m> upstream=<id> TIMEOUT: первый байт не пришёл за <N>ms → фейловер`
(+ `(backoff: отдых …)` если задан backoff), метрика
`llm_proxy_requests_total{result="timeout"}`;
- закрытие апстримом пустого стрима (EOF без единого байта) — **не** таймаут:
это валидный короткий ответ, он проксируется как обычно.
Значение — ISO-8601-длительность, формат как у `backoff` (`PT30S`,
`PT1M30S`, `P1D`, …).
```yaml
providers:
- id: routerai
url: "https://routerai.ru/api/v1"
first_byte_timeout: PT30S # окно для всех моделей провайдера
upstreams:
- id: routerai-gpt4o
provider: routerai
model: gpt-4o
first_byte_timeout: PT0S # у модели watchdog отключён
```
При старте выводится, что настроено:
`[llm-proxy] first_byte_timeout: upstreams=routerai-gpt4o=PT0S providers=routerai=PT30S default=PT30S`.
### Фоновый пинг апстримов в откате (`probe_interval`)
Когда апстрим «уходит в аут» надолго (backoff растёт до потолка), без пинга
он вернётся в ротацию только по истечении интервала отката — это может быть
минуты простоя, хотя модель уже могла ожить. `probe_interval` включает
фоновый «пингер»: пока апстрим в backoff-откате, HealthProber каждые
`probe_interval` шлёт ему минимальный запрос (`2+2=?`, `max_tokens=4`,
`stream=true`) и слушает первый байт ответа. Как только апстрим ответил —
его снимают с отката немедленно (`backoff.recordSuccess`), и он возвращается
в общий пул, не дожидаясь истечения backoff-интервала.
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].probe_interval` | ISO-8601-длительность / отсутствует | Интервал пинга для всех моделей провайдера |
| `upstreams[].probe_interval` | ISO-8601-длительность / отсутствует | Интервал для конкретной модели (приоритет над провайдерским) |
**Дефолт** (если не задан нигде): `PT1S`.
**Приоритет — модели.** Если `upstreams[].probe_interval` задан, он
перекрывает провайдерский; если нет — берётся провайдерский; если нет и там —
дефолт `PT1S`.
**Отключение:** `PT0S` на провайдере или апстриме отключает пинг для этого
апстрима — тогда он возвращается в пул только по истечении backoff-интервала.
Особенности:
- пингуются **только** апстримы, которые сейчас в backoff-откате; живой
апстрим пинг не получает (проверка состояния — каждые `probe_interval`);
- пинг **не занимает слот конкурентности** (идёт в обход `max_concurrency`),
чтобы не отжимать слот у живых запросов;
- ошибка пинга **не** считается в backoff (иначе каждая попытка удваивала
бы интервал отката); учитывается только успех;
- признак «ответил» — HTTP 2xx и первый байт тела за `first_byte_timeout`
(тот же watchdog, что и у обычных запросов);
- на пинг применяются патчи провайдера и апстрима (как к обычному запросу),
чтобы он валидировал реальный путь;
- при восстановлении — лог:
`[llm-proxy] probe upstream=<id> ответил — снят с backoff, вернулся в пул`,
метрика `llm_proxy_probe_pings_total{result="ok"}`;
- при молчании — лог `[llm-proxy] probe upstream=<id> молчит — остаюсь в backoff`,
метрика `llm_proxy_probe_pings_total{result="fail"}`.
Значение — ISO-8601-длительность, формат как у `backoff` (`PT1S`, `PT500MS`,
`PT5S`, …).
```yaml
providers:
- id: routerai
url: "https://routerai.ru/api/v1"
probe_interval: PT1S # пингуем все модели провайдера раз в секунду
upstreams:
- id: routerai-gpt4o
provider: routerai
model: gpt-4o
probe_interval: PT0S # у модели пинг отключён (ждать истечения backoff)
```
При старте выводится, что настроено:
`[llm-proxy] probe_interval: upstreams=routerai-gpt4o=PT0S providers=routerai=PT1S default=PT1S`.
### Prometheus-метрики (`/metrics`)
Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics`
@@ -331,7 +460,8 @@ upstreams:
| Метрика | Тип | Смысл |
|---|---|---|
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `cancelled` — клиент отвалился посреди стрима |
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `timeout` — апстрим не прислал первый байт стрима за `first_byte_timeout`, `cancelled` — клиент отвалился посреди стрима |
| `llm_proxy_probe_pings_total{upstream, provider, result}` | counter | фоновые пинги апстримов в backoff-откате (HealthProber). `result`: `ok` — апстрим ответил (снят с отката), `fail` — молчит (остался в откате) |
| `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 = жив) |
@@ -343,7 +473,7 @@ upstreams:
- оборот по виртуальным моделям: `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)`.
- доля ошибок провайдера: `sum(rate(llm_proxy_requests_total{result=~"4xx|429|402|5xx|net_err|timeout"}[10m])) by (provider) / sum(rate(llm_proxy_requests_total[10m])) by (provider)`.
### Пример сборки тела (многослойный `patch`)
@@ -429,10 +559,11 @@ data class ProviderConf(
val think_tags: String? = null,
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
val first_byte_timeout: Duration? = null, // окно первого байта стрима (ISO-8601)
val probe_interval: Duration? = null, // интервал фонового пинга (ISO-8601)
)
@Serializable
data class UpstreamConf(
val id: String, // наш внутренний id апстрима
val provider: String, // ссылка на ProviderConf.id
@@ -441,6 +572,8 @@ data class UpstreamConf(
val patch: JsonObject? = null, // к запросам этой апстрим-модели
val think_tags: String? = null, // переопределение think-режима модели
val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским)
val first_byte_timeout: Duration? = null, // окно первого байта стрима (приоритет над провайдерским)
val probe_interval: Duration? = null, // интервал фонового пинга (приоритет над провайдерским)
)
@Serializable
@@ -605,6 +738,10 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
(`(backoff: отдых PT…)`); при пропуске апстрима в откате —
`upstream=<id> в backoff-откате (~Ns) — пропускаю`; при старте — список
настроенных backoff (`backoff: upstreams=… providers=…`).
- **Таймаут первого байта стрима**: при срабатывании watchdog'а —
`upstream=<id> TIMEOUT: первый байт не пришёл за <N>ms → фейловер`;
при старте — список настроенных `first_byte_timeout`
(`first_byte_timeout: upstreams=… providers=… default=PT30S`).
- **Все в откате** — `503 all upstreams cooling, retry in Ns` для
`model=<витрина>` с заголовком `Retry-After: N`.
+14
View File
@@ -30,6 +30,20 @@ non-stream.
1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате
— `503` с `Retry-After`.
Для `stream=true` запросов действует watchdog на первый байт ответа:
`first_byte_timeout` (ISO-8601, дефолт `PT30S`, `PT0S` отключает) — если
апстрим не прислал ни одного байта за окно, слот освобождается, ошибка идёт
в backoff и роутер фейловерит на следующего свободного апстрима. Приоритет:
`upstreams[].first_byte_timeout` → `providers[].first_byte_timeout` → `PT30S`.
Пока апстрим «отдыхает» в backoff-откате, его можно фоново пинговать:
`probe_interval` (ISO-8601, дефолт `PT1S`, `PT0S` отключает) — HealthProber
шлёт апстриму минимальный запрос (`2+2=?`, `max_tokens=4`) и, как только тот
ответил, снимает его с отката немедленно, не дожидаясь истечения
backoff-интервала. Пинг не занимает слот конкурентности и не засчитывается
в backoff. Приоритет: `upstreams[].probe_interval` →
`providers[].probe_interval` → `PT1S`.
Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по
`model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) —
подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список
@@ -0,0 +1,133 @@
package pw.binom.llmproxy
import io.ktor.client.HttpClient
import io.ktor.client.call.body
import io.ktor.client.request.headers
import io.ktor.client.request.preparePost
import io.ktor.client.request.setBody
import io.ktor.utils.io.ByteReadChannel
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive
import kotlin.time.Duration
/**
* Фоновый «пингер» апстримов, которые в данный момент в backoff-откате.
*
* Пока апстрим молчит (в откате), HealthProber раз в [interval] шлёт ему
* минимальный chat.completion (`2+2=?`, `max_tokens=4`) и слушает первый байт
* ответа. Как только апстрим ответил — вызывается [BackoffRegistry.recordSuccess],
* и он возвращается в общий пул сразу, не дожидаясь истечения backoff-интервала.
* Это ускоряет восстановление после долгого простоя провайдера.
*
* Пинг НЕ занимает слот конкурентности ([UpstreamCounter]) — он идёт в обход,
* чтобы не отжимать слот у живых запросов. Ошибки пинга НЕ считаются в backoff
* (иначе каждая попытка удваивала бы интервал отката); считается только успех.
*/
class HealthProber(
private val providersById: Map<String, ProviderConf>,
private val upstreams: List<UpstreamConf>,
private val http: HttpClient,
private val backoff: BackoffRegistry,
private val metrics: MetricsRegistry,
private val scope: CoroutineScope,
) {
private val jobs = mutableMapOf<String, Job>()
/** Запустить пинг-корутины для всех апстримов с включённым пингом. */
fun start() {
for (up in upstreams) {
val interval = effectiveProbeInterval(up, providersById[up.provider])
if (interval != null) {
jobs[up.id] = scope.launch { probeLoop(up, interval) }
}
}
}
/** Цикл пинга одного апстрима: пингуем только пока он в откате. */
internal suspend fun probeLoop(up: UpstreamConf, interval: Duration, pingImpl: suspend (UpstreamConf) -> Boolean = ::ping) {
while (currentCoroutineContext().isActive) {
if (backoff.isCooling(up)) {
val healthy = pingImpl(up)
metrics.recordProbe(up, if (healthy) "ok" else "fail")
if (healthy) {
backoff.recordSuccess(up)
log.info { "[llm-proxy] probe upstream=${up.id} ответил — снят с backoff, вернулся в пул" }
} else {
log.info { "[llm-proxy] probe upstream=${up.id} молчит — остаюсь в backoff" }
}
}
delay(interval)
}
}
/** Одна пинг-попытка: минимальный stream-запрос, слушаем первый байт. */
private suspend fun ping(up: UpstreamConf): Boolean {
val provider = providersById[up.provider] ?: return false
val url = provider.url.trimEnd('/') + "/chat/completions"
val providerKey = resolveEnv(provider.key)
val body = buildProbeBody(up, provider)
val firstByteTimeout = effectiveFirstByteTimeout(up, provider) ?: DEFAULT_FIRST_BYTE_TIMEOUT
return try {
http.preparePost(url) {
headers {
if (providerKey.isNotEmpty()) {
set("Authorization", "Bearer $providerKey")
}
set("Content-Type", "application/json")
}
setBody(body.toString())
}.execute { resp ->
if (resp.status.value !in 200..299) {
return@execute false
}
val ch = resp.body<ByteReadChannel>()
val buf = ByteArray(64)
readAvailableWithTimeout(ch, buf, firstByteTimeout.inWholeMilliseconds) > 0
}
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
log.info { "[llm-proxy] probe upstream=${up.id} ошибка: ${e.message}" }
false
}
}
}
/**
* Собрать тело пинг-запроса: минимальный chat.completion, на который любой
* живой апстрим ответит сразу. Применяются патчи провайдера и апстрима (в том
* же порядке, что и для реальных запросов), чтобы пинг валидировал реальный
* путь запроса, но без клиентских данных и модели.
*/
internal fun buildProbeBody(up: UpstreamConf, provider: ProviderConf?): JsonObject {
val base = JsonObject(
mapOf(
"model" to JsonPrimitive(up.model),
"messages" to JsonArray(
listOf(
JsonObject(
mapOf(
"role" to JsonPrimitive("user"),
"content" to JsonPrimitive("2+2=?"),
),
),
),
),
"max_tokens" to JsonPrimitive(4),
"stream" to JsonPrimitive(true),
),
)
var acc = base
listOf(provider?.patch, up.patch).filterNotNull().forEach { patch ->
acc = merge(acc, patch)
}
return acc
}
+153 -6
View File
@@ -28,6 +28,14 @@ import io.ktor.utils.io.ByteReadChannel
import io.ktor.utils.io.ByteWriteChannel
import io.ktor.utils.io.writeFully
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.async
import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.selects.onTimeout
import kotlinx.coroutines.selects.select
import kotlinx.coroutines.sync.Mutex
import kotlinx.datetime.Clock
import kotlinx.serialization.Serializable
@@ -103,6 +111,20 @@ fun main() {
" providers=" +
config.providers.filter { it.backoff != null }.map { p -> "${p.id}=${p.backoff?.toIsoString()}" }.joinToString(", ")
}
log.info {
"[llm-proxy] first_byte_timeout: upstreams=" +
config.upstreams.filter { it.first_byte_timeout != null }.map { up -> "${up.id}=${up.first_byte_timeout?.toIsoString()}" }.joinToString(", ") +
" providers=" +
config.providers.filter { it.first_byte_timeout != null }.map { p -> "${p.id}=${p.first_byte_timeout?.toIsoString()}" }.joinToString(", ") +
" default=${DEFAULT_FIRST_BYTE_TIMEOUT.toIsoString()}"
}
log.info {
"[llm-proxy] probe_interval: upstreams=" +
config.upstreams.filter { it.probe_interval != null }.map { up -> "${up.id}=${up.probe_interval?.toIsoString()}" }.joinToString(", ") +
" providers=" +
config.providers.filter { it.probe_interval != null }.map { p -> "${p.id}=${p.probe_interval?.toIsoString()}" }.joinToString(", ") +
" default=${DEFAULT_PROBE_INTERVAL.toIsoString()}"
}
log.info {
"[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}"
}
@@ -112,6 +134,8 @@ fun main() {
val http = createHttpClient()
val sessions = SessionRegistry()
val prober = HealthProber(providersById, config.upstreams, http, backoff, metrics, CoroutineScope(SupervisorJob() + Dispatchers.Default))
prober.start()
startServer(config.server.host, config.server.port) {
proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
}
@@ -279,13 +303,14 @@ private suspend fun handleChat(
val ct = resp.headers["Content-Type"] ?: "text/event-stream"
val status = HttpStatusCode.fromValue(upstreamStatus)
val thinkMode = effectiveThinkTags(up, provider)
val firstByteTimeout = effectiveFirstByteTimeout(up, provider)
call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>()
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
streamRawWithDoneContract(ch)
streamRawWithDoneContract(ch, firstByteTimeout)
} else {
streamSseWithThinkTags(ch, thinkMode)
streamSseWithThinkTags(ch, thinkMode, firstByteTimeout)
}
}
metrics.record(modelName, up, "ok")
@@ -322,6 +347,15 @@ private suspend fun handleChat(
if (failover) continue
if (responded) return
} catch (e: FirstByteTimeoutException) {
val rest = backoff.recordFailure(up)
metrics.record(modelName, up, "timeout")
log.warn {
"[llm-proxy] chat model=$modelName upstream=${up.id} TIMEOUT: первый байт не пришёл за ${e.timeoutMs}ms → фейловер" +
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
}
failed.add(up.id)
continue
} catch (e: CancellationException) {
metrics.record(modelName, up, "cancelled")
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
@@ -493,6 +527,38 @@ internal fun effectiveThinkTags(up: UpstreamConf, provider: ProviderConf?): Stri
}
}
/**
* Эффективный first-byte-timeout для стрим-ответа: значение у апстрима, если
* задано; иначе у провайдера; иначе [DEFAULT_FIRST_BYTE_TIMEOUT]. Возвращает
* null, если явно отключено (`Duration.ZERO` — `PT0S`): тогда стрим
* ждёт первого байта без лимита (в пределах общего
* [UPSTREAM_REQUEST_TIMEOUT_MS]).
*/
internal fun effectiveFirstByteTimeout(up: UpstreamConf, provider: ProviderConf?): Duration? {
val raw = up.first_byte_timeout ?: provider?.first_byte_timeout ?: DEFAULT_FIRST_BYTE_TIMEOUT
return if (raw <= Duration.ZERO) null else raw
}
/**
* Эффективный интервал фонового пинга апстрима в backoff-откате: значение у
* апстрима, если задано; иначе у провайдера; иначе [DEFAULT_PROBE_INTERVAL].
* Возвращает null, если явно отключено (`Duration.ZERO` — `PT0S`): тогда
* HealthProber этот апстрим не пингует (ждём естественного истечения бэкоффа).
*/
internal fun effectiveProbeInterval(up: UpstreamConf, provider: ProviderConf?): Duration? {
val raw = up.probe_interval ?: provider?.probe_interval ?: DEFAULT_PROBE_INTERVAL
return if (raw <= Duration.ZERO) null else raw
}
/**
* Сигнализирует, что апстрим не прислал ни одного байта тела стрим-ответа за
* [timeoutMs] мс. Обрабатывается как обычный сбой (фейловер, backoff,
* метрика `timeout`), но отделён от [CancellationException] (отмена клиентом)
* и от generic `Exception` (сетевые ошибки).
*/
class FirstByteTimeoutException(val timeoutMs: Long) :
RuntimeException("upstream did not send first byte within ${timeoutMs}ms")
/** Освободить слот апстрима (в finally по завершении проксирования). */
internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) {
active.getValue(up.id).release()
@@ -827,12 +893,26 @@ internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDO
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
* Если задан [firstByteTimeout] (non-null), первый байт тела должен прийти
* за это время, иначе бросаем [FirstByteTimeoutException]; после первого
* байта watchdog выключается (ждём ответа вплоть до общего
* [UPSTREAM_REQUEST_TIMEOUT_MS]).
*/
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(
source: ByteReadChannel,
firstByteTimeout: Duration? = null,
) {
val detector = StreamEndDetector()
val buf = ByteArray(8192)
val timeoutMs = firstByteTimeout?.inWholeMilliseconds ?: 0L
var isFirstRead = timeoutMs > 0L
while (true) {
val n = source.readAvailable(buf)
val n = if (isFirstRead) {
isFirstRead = false
readAvailableWithTimeout(source, buf, timeoutMs)
} else {
source.readAvailable(buf)
}
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
@@ -853,14 +933,30 @@ internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteRead
* записывается как `data: <json>\n\n` (событие-граница SSE), а при ошибке парса —
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
* Если задан [firstByteTimeout] (non-null), первая строка должна прийти за это
* время, иначе [FirstByteTimeoutException] (watchdog только на первой строке).
*/
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(
source: ByteReadChannel,
thinkMode: String,
firstByteTimeout: Duration? = null,
) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
val timeoutMs = firstByteTimeout?.inWholeMilliseconds ?: 0L
var isFirstRead = timeoutMs > 0L
suspend fun readNext(): String? = if (isFirstRead) {
isFirstRead = false
readLineWithTimeout(source, timeoutMs)
} else {
source.readLine(LineEnding.Lenient)
}
while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break
val line = readNext() ?: break
when {
line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim()
@@ -913,6 +1009,49 @@ internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadCha
}
}
/**
* Прочитать из [source] в [buf] с жёстким лимитом [timeoutMs]: если за это
* время ни одного байта не пришло, бросить [FirstByteTimeoutException].
* `select` отменяет «проигравшую» ветку; при таймауте явно отменяем и
* дожидаемся чтение (`cancelAndJoin`), чтобы соединение с апстримом
* корректно закрылось. Сентинел [Int.MIN_VALUE] невозможен как результат
* [ByteReadChannel.readAvailable] (это `-1` для EOF или `>= 0` для байт).
*/
internal suspend fun readAvailableWithTimeout(source: ByteReadChannel, buf: ByteArray, timeoutMs: Long): Int =
coroutineScope {
val readJob = async { source.readAvailable(buf) }
val result = select<Int> {
readJob.onAwait { it }
onTimeout(timeoutMs) { Int.MIN_VALUE }
}
if (result == Int.MIN_VALUE) {
readJob.cancelAndJoin()
throw FirstByteTimeoutException(timeoutMs)
}
result
}
/**
* Прочитать строку из [source] с лимитом [timeoutMs]: если за это время ни
* одной строки/байта не пришло, бросить [FirstByteTimeoutException].
* [ByteReadChannel.readLine] возвращает `null` и при EOF, и (через `onTimeout`)
* при таймауте; различаем по состоянию [readJob]: если ветка чтения не
* завершилась к моменту возврата из `select` — таймаут (select отменил её).
*/
internal suspend fun readLineWithTimeout(source: ByteReadChannel, timeoutMs: Long): String? =
coroutineScope {
val readJob = async { source.readLine(LineEnding.Lenient) }
val result = select<String?> {
readJob.onAwait { it }
onTimeout(timeoutMs) { null }
}
if (result == null && !readJob.isCompleted) {
readJob.cancelAndJoin()
throw FirstByteTimeoutException(timeoutMs)
}
result
}
/**
* Пересобрать SSE-чанк: по каждому choice (ключ `index`, дефолт 0) взять
* `delta.content` (только если это JSON-строка; массив частей не трогаем) и
@@ -966,6 +1105,8 @@ data class ProviderConf(
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
val probe_interval: Duration? = null,
)
data class UpstreamConf(
@@ -976,6 +1117,8 @@ data class UpstreamConf(
val patch: JsonObject? = null,
val think_tags: String? = null,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
val probe_interval: Duration? = null,
)
data class ModelConf(
@@ -1034,6 +1177,8 @@ internal fun parseConfig(root: YamlElement): Config {
reasoning_field = m.strOrNull("reasoning_field"),
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
backoff = m.durationOrNull("backoff"),
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
probe_interval = m.durationOrNull("probe_interval"),
)
}
@@ -1047,6 +1192,8 @@ internal fun parseConfig(root: YamlElement): Config {
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
think_tags = m.strOrNull("think_tags"),
backoff = m.durationOrNull("backoff"),
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
probe_interval = m.durationOrNull("probe_interval"),
)
}
@@ -9,7 +9,9 @@ import kotlinx.coroutines.sync.withLock
*
* Метрики (всё в памяти, сброс при рестарте):
* - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | cancelled`;
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | timeout | cancelled`;
* - `llm_proxy_probe_pings_total{upstream,provider,result}` — счётчик фоновых
* пингов апстримов в backoff-откате (HealthProber); `result`: `ok | fail`;
* - `llm_proxy_upstream_inflight{upstream,provider}` — сейчас в работе (слоты);
* - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев
* (счётчик backoff);
@@ -19,6 +21,7 @@ import kotlinx.coroutines.sync.withLock
class MetricsRegistry {
private val lock = Mutex()
private val requests = HashMap<RequestKey, Long>()
private val probes = HashMap<ProbeKey, Long>()
private data class RequestKey(
val model: String,
@@ -27,6 +30,12 @@ class MetricsRegistry {
val result: String,
)
private data class ProbeKey(
val upstream: String,
val provider: String,
val result: String,
)
/** Учёт chat-запроса по исходу (значения `result` — в описании класса). */
suspend fun record(model: String, up: UpstreamConf, result: String) {
lock.withLock {
@@ -35,6 +44,14 @@ class MetricsRegistry {
}
}
/** Учёт пинга апстрима в backoff-откате (result: `ok`/`fail`). */
suspend fun recordProbe(up: UpstreamConf, result: String) {
lock.withLock {
val key = ProbeKey(up.id, up.provider, result)
probes[key] = (probes[key] ?: 0L) + 1
}
}
/** Текстовая выгрузка в формате Prometheus на текущий момент. */
suspend fun render(
config: Config,
@@ -56,6 +73,19 @@ class MetricsRegistry {
)
}
}
sb.append("# HELP llm_proxy_probe_pings_total Фоновые пинги апстримов в backoff-откате\n")
sb.append("# TYPE llm_proxy_probe_pings_total counter\n")
lock.withLock {
probes.entries
.sortedWith(compareBy({ it.key.upstream }, { it.key.provider }, { it.key.result }))
.forEach { (k, n) ->
sb.append(
"llm_proxy_probe_pings_total{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) {
@@ -1,4 +1,26 @@
package pw.binom.llmproxy
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
/** Предел времени на весь запрос к апстриму, мс: от отправки до приёма всего ответа (включая стриминг). */
const val UPSTREAM_REQUEST_TIMEOUT_MS: Long = 5 * 60_000L
/**
* Default «watchdog» на первый байт тела стрим-ответа: если апстрим не прислал
* ни одного байта за это время, прокси считает это сбоем апстрима и
* фейловерит (освобождает слот конкурентности, считает ошибку в backoff и
* идёт к следующему апстриму модели). Применяется ТОЛЬКО к stream=true
* запросам; для non-stream работает общий [UPSTREAM_REQUEST_TIMEOUT_MS].
* Отключается явным `first_byte_timeout: PT0S` на апстриме или провайдере.
*/
val DEFAULT_FIRST_BYTE_TIMEOUT: Duration = 30.seconds
/**
* Default интервал фоновой проверки («пинга») апстрима, который сейчас в
* backoff-откате: пока апстрим молчит, HealthProber шлёт ему минимальный
* запрос (2+2=?, max_tokens=4) каждые [DEFAULT_PROBE_INTERVAL], и как только
* он отвечает — снимает его с отката раньше, чем истечёт интервал бэкоффа.
* Отключается явным `probe_interval: PT0S` на апстриме или провайдере.
*/
val DEFAULT_PROBE_INTERVAL: Duration = 1.seconds
@@ -4,6 +4,7 @@ import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.seconds
import io.ktor.http.headersOf
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.Json
@@ -586,4 +587,60 @@ class ConfigLogicTest {
assertEquals(null, cfg.providers[0].reasoning_field)
assertEquals(false, cfg.providers[0].reasoning_empty_ok)
}
@Test
fun parseConfigReadsFirstByteTimeoutOnProviderAndUpstream() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
first_byte_timeout: PT45S
- id: p2
url: "https://y.ru/api/v1"
upstreams:
- id: u1
provider: p1
model: real-1
first_byte_timeout: PT0S
- id: u2
provider: p2
model: real-2
models:
- name: m1
upstreams: [u1, u2]
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals(45.seconds, cfg.providers[0].first_byte_timeout)
assertEquals(null, cfg.providers[1].first_byte_timeout)
assertEquals(0.seconds, cfg.upstreams[0].first_byte_timeout)
assertEquals(null, cfg.upstreams[1].first_byte_timeout)
}
@Test
fun parseConfigReadsProbeIntervalOnProviderAndUpstream() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
probe_interval: PT2S
- id: p2
url: "https://y.ru/api/v1"
upstreams:
- id: u1
provider: p1
model: real-1
probe_interval: PT0S
- id: u2
provider: p2
model: real-2
models:
- name: m1
upstreams: [u1, u2]
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals(2.seconds, cfg.providers[0].probe_interval)
assertEquals(null, cfg.providers[1].probe_interval)
assertEquals(0.seconds, cfg.upstreams[0].probe_interval)
assertEquals(null, cfg.upstreams[1].probe_interval)
}
}
@@ -0,0 +1,218 @@
package pw.binom.llmproxy
import io.ktor.utils.io.ByteChannel
import io.ktor.utils.io.close
import io.ktor.utils.io.readAvailable
import io.ktor.utils.io.writeFully
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.milliseconds
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.advanceTimeBy
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
/**
* Тесты watchdog'а на первый байт тела стрим-ответа.
*
* Сценарии:
* - Апстрим не прислал ничего за N мс → [FirstByteTimeoutException] на первой
* попытке чтения (после первой попытки watchdog отключается).
* - Апстрим закрыл канал до таймаута, не прислав данных → возврат EOF (-1 /
* null), НЕ [FirstByteTimeoutException].
* - Апстрим прислал первый байт вовремя → читаем его, watchdog отключён.
* - Когда таймаут не задан (`null`) → читаем без лимита (как и было).
*/
class FirstByteTimeoutTest {
/** Окно для watchdog'а в тестах: виртуальное время `runTest` его мгновенно прокручивает. */
private val watchWindow = 100.milliseconds
@Test
fun readAvailableTimesOutWhenNoBytesAndNotClosed() = runTest {
val src = ByteChannel(autoFlush = true)
val buf = ByteArray(64)
val done = CompletableDeferred<Throwable?>()
launch {
try {
readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
assertEquals(watchWindow.inWholeMilliseconds, e.timeoutMs)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
}
@Test
fun readAvailableReturnsEofImmediatelyWhenChannelClosedEmpty() = runTest {
// Канал сразу закрыт без данных: -1 (EOF), никакого таймаута.
val src = ByteChannel(autoFlush = true)
src.close(null)
val buf = ByteArray(64)
val n = readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
assertEquals(-1, n, "закрытый пустой канал → EOF, не таймаут")
}
@Test
fun readAvailableReturnsBytesWhenDataArrivesInTime() = runTest {
val src = ByteChannel(autoFlush = true)
val payload = "hello".encodeToByteArray()
src.writeFully(payload)
val buf = ByteArray(64)
val n = readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
assertEquals(payload.size, n)
assertEquals(payload.decodeToString(), buf.decodeToString(0, n))
}
@Test
fun readLineTimesOutWhenNoLinesAndNotClosed() = runTest {
val src = ByteChannel(autoFlush = true)
val done = CompletableDeferred<Throwable?>()
launch {
try {
readLineWithTimeout(src, watchWindow.inWholeMilliseconds)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
}
@Test
fun readLineReturnsNullWhenChannelClosedEmpty() = runTest {
// Закрытый пустой канал → null (EOF), не таймаут. Это критично: иначе бы
// мы фейловерили легитимные короткие ответы, где апстрим закрыл стрим
// до первой строки.
val src = ByteChannel(autoFlush = true)
src.close(null)
val line = readLineWithTimeout(src, watchWindow.inWholeMilliseconds)
assertEquals(null, line)
}
@Test
fun rawStreamTimesOutWhenUpstreamSilent() = runTest {
// Сырой passthrough: апстрим молчит вечно → на первой попытке чтения
// вылетает FirstByteTimeoutException, в out ничего не пишется.
val src = ByteChannel(autoFlush = true)
val out = ByteChannel(autoFlush = true)
val done = CompletableDeferred<Throwable?>()
launch {
try {
out.streamRawWithDoneContract(src, watchWindow)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
// Подтверждаем, что в out ничего не утекло.
out.close(null)
val buf = ByteArray(64)
val n = out.readAvailable(buf)
assertEquals(-1, n, "до таймаута клиенту ничего не должно уйти")
}
@Test
fun thinkStreamTimesOutWhenUpstreamSilent() = runTest {
// SSE-вариант: первая строка не приходит → FirstByteTimeoutException,
// в out ничего не пишется.
val src = ByteChannel(autoFlush = true)
val out = ByteChannel(autoFlush = true)
val done = CompletableDeferred<Throwable?>()
launch {
try {
out.streamSseWithThinkTags(src, "split", watchWindow)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
out.close(null)
val buf = ByteArray(64)
val n = out.readAvailable(buf)
assertEquals(-1, n)
}
@Test
fun thinkStreamFirstLineInTimeThenReadsUnlimited() = runTest {
// Первый байт/строка пришли вовремя → watchdog выключен, остальной
// поток читается обычным порядком. Проверяем, что НЕ бросило
// FirstByteTimeoutException и данные дошли до out.
val src = ByteChannel(autoFlush = true)
val out = ByteChannel(autoFlush = true)
val payload = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
src.writeFully(payload.encodeToByteArray())
src.close(null)
out.streamSseWithThinkTags(src, "split", watchWindow)
out.close(null)
val sb = StringBuilder()
val buf = ByteArray(256)
while (true) {
val n = out.readAvailable(buf)
if (n == -1) break
if (n > 0) sb.append(buf.decodeToString(0, n))
}
assertTrue(sb.toString().contains("data:"), "данные дошли до out: $sb")
}
@Test
fun effectiveFirstByteTimeoutPrefersUpstreamThenProviderThenDefault() {
// upstream > provider > default
val provider = ProviderConf("p", "https://x", "", null, null, first_byte_timeout = 20.milliseconds)
val upstream = UpstreamConf("u", "p", "m", null, null, null, null, first_byte_timeout = 5.milliseconds)
val upstreamNoOwn = UpstreamConf("u", "p", "m")
assertEquals(5.milliseconds, effectiveFirstByteTimeout(upstream, provider))
assertEquals(20.milliseconds, effectiveFirstByteTimeout(upstreamNoOwn, provider))
assertEquals(DEFAULT_FIRST_BYTE_TIMEOUT, effectiveFirstByteTimeout(upstreamNoOwn, null))
}
@Test
fun effectiveFirstByteTimeoutZeroDisablesWatchdog() {
// 0s / PT0S → null: passthrough без watchdog'а.
val up = UpstreamConf("u", "p", "m", null, null, null, null, first_byte_timeout = 0.milliseconds)
assertEquals(null, effectiveFirstByteTimeout(up, null))
val providerZero = ProviderConf("p", "https://x", "", null, null, first_byte_timeout = 0.milliseconds)
val upNoOwn = UpstreamConf("u", "p", "m")
assertEquals(null, effectiveFirstByteTimeout(upNoOwn, providerZero))
}
}
@@ -0,0 +1,148 @@
package pw.binom.llmproxy
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.TestScope
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.minutes
import kotlin.time.Duration.Companion.seconds
class HealthProberTest {
private val json = Json { ignoreUnknownKeys = true }
@Test
fun effectiveProbeIntervalPrefersUpstreamThenProviderThenDefault() {
val provider = ProviderConf("p", "https://x", "", null, null, probe_interval = 2.seconds)
val upstream = UpstreamConf("u", "p", "m", null, null, null, null, null, probe_interval = 5.milliseconds)
val upstreamNoOwn = UpstreamConf("u2", "p", "m")
assertEquals(5.milliseconds, effectiveProbeInterval(upstream, provider))
assertEquals(2.seconds, effectiveProbeInterval(upstreamNoOwn, provider))
assertEquals(DEFAULT_PROBE_INTERVAL, effectiveProbeInterval(upstreamNoOwn, null))
}
@Test
fun effectiveProbeIntervalZeroDisablesProber() {
val up = UpstreamConf("u", "p", "m", null, null, null, null, null, probe_interval = 0.milliseconds)
assertNull(effectiveProbeInterval(up, null))
val providerZero = ProviderConf("p", "https://x", "", null, null, probe_interval = 0.milliseconds)
assertNull(effectiveProbeInterval(UpstreamConf("u2", "p", "m"), providerZero))
}
@Test
fun buildProbeBodyIsMinimalChatWithModelReplacedAndPatchesApplied() {
val provider = ProviderConf(
"p", "https://x", "", null,
Json.parseToJsonElement("""{"provider":{"allow_fallbacks":false}}""").jsonObject,
probe_interval = 1.seconds,
)
val up = UpstreamConf(
"u", "p", "real-1", null,
Json.parseToJsonElement("""{"temperature":0.5}""").jsonObject,
null, null, null, 1.seconds,
)
val body = buildProbeBody(up, provider)
assertEquals("real-1", body["model"]?.jsonPrimitive?.content)
assertEquals("2+2=?", body["messages"]?.jsonArray?.first()?.jsonObject?.get("content")?.jsonPrimitive?.content)
assertEquals("true", body["stream"]?.jsonPrimitive?.content)
assertEquals("4", body["max_tokens"]?.jsonPrimitive?.content)
assertEquals("false", body["provider"]?.jsonObject?.get("allow_fallbacks")?.jsonPrimitive?.content)
assertEquals("0.5", body["temperature"]?.jsonPrimitive?.content)
}
@Test
fun buildProbeBodyIgnoresModelPatch() {
val up = UpstreamConf("u", "p", "real-1")
val provider = ProviderConf("p", "https://x")
val body = buildProbeBody(up, provider)
assertEquals(
"""{"model":"real-1","messages":[{"role":"user","content":"2+2=?"}],"max_tokens":4,"stream":true}""",
body.toString(),
)
}
@Test
fun probeLoopRecoversUpstreamWhenPingSucceeds() = runTest {
val provider = ProviderConf("p", "https://x", "", null, null, backoff = 5.minutes)
val up = UpstreamConf("u", "p", "m", null, null, null, backoff = 5.minutes, probe_interval = 5.milliseconds)
val backoff = BackoffRegistry(mapOf("p" to provider), mapOf("u" to up))
backoff.recordFailure(up)
assertEquals(true, backoff.isCooling(up))
val metrics = MetricsRegistry()
val prober = HealthProber(
providersById = mapOf("p" to provider),
upstreams = listOf(up),
http = createHttpClient(),
backoff = backoff,
metrics = metrics,
scope = this,
)
var pings = 0
val job = launch {
prober.probeLoop(up, 5.milliseconds) { pinged ->
pings++
assertEquals("u", pinged.id)
true
}
}
runUntil { pings >= 1 }
// успех уже зафиксирован (recordSuccess) — даём циклу несколько
// итераций и убеждаемся, что повторных пингов не было
delay(10.milliseconds)
job.cancel()
assertEquals(1, pings)
assertEquals(false, backoff.isCooling(up))
}
@Test
fun probeLoopKeepsCoolingWhenPingFails() = runTest {
val provider = ProviderConf("p", "https://x", "", null, null, backoff = 5.minutes)
val up = UpstreamConf("u", "p", "m", null, null, null, backoff = 5.minutes, probe_interval = 5.milliseconds)
val backoff = BackoffRegistry(mapOf("p" to provider), mapOf("u" to up))
backoff.recordFailure(up)
val metrics = MetricsRegistry()
val prober = HealthProber(
providersById = mapOf("p" to provider),
upstreams = listOf(up),
http = createHttpClient(),
backoff = backoff,
metrics = metrics,
scope = this,
)
var pings = 0
val job = launch {
prober.probeLoop(up, 5.milliseconds) { _ ->
pings++
false
}
}
runUntil { pings >= 2 }
job.cancel()
assertEquals(true, backoff.isCooling(up))
}
private suspend fun TestScope.runUntil(condition: () -> Boolean) {
var i = 0
while (!condition() && i < 10_000) {
delay(1)
i++
}
if (!condition()) throw AssertionError("condition not met in time")
}
}