feat: фоновый пинг апстримов в backoff-откате (probe_interval) — мгновенное возвращение в пул
Build LLM Proxy / Build and push (release) Successful in 42s

This commit is contained in:
2026-09-15 01:00:10 +03:00
parent 249b83eb51
commit c6a24a94b0
8 changed files with 452 additions and 0 deletions
+69
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
@@ -385,6 +389,68 @@ upstreams:
При старте выводится, что настроено:
`[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`
@@ -395,6 +461,7 @@ upstreams:
| Метрика | Тип | Смысл |
|---|---|---|
| `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 = жив) |
@@ -494,6 +561,7 @@ data class ProviderConf(
val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
val first_byte_timeout: Duration? = null, // окно первого байта стрима (ISO-8601)
val probe_interval: Duration? = null, // интервал фонового пинга (ISO-8601)
)
data class UpstreamConf(
@@ -505,6 +573,7 @@ data class UpstreamConf(
val think_tags: String? = null, // переопределение think-режима модели
val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским)
val first_byte_timeout: Duration? = null, // окно первого байта стрима (приоритет над провайдерским)
val probe_interval: Duration? = null, // интервал фонового пинга (приоритет над провайдерским)
)
@Serializable
+8
View File
@@ -36,6 +36,14 @@ non-stream.
в 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
}
@@ -28,6 +28,9 @@ 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
@@ -115,6 +118,13 @@ fun main() {
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}"
}
@@ -124,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)
}
@@ -527,6 +539,17 @@ internal fun effectiveFirstByteTimeout(up: UpstreamConf, provider: ProviderConf?
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,
@@ -1083,6 +1106,7 @@ data class ProviderConf(
val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
val probe_interval: Duration? = null,
)
data class UpstreamConf(
@@ -1094,6 +1118,7 @@ data class UpstreamConf(
val think_tags: String? = null,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
val probe_interval: Duration? = null,
)
data class ModelConf(
@@ -1153,6 +1178,7 @@ internal fun parseConfig(root: YamlElement): Config {
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"),
)
}
@@ -1167,6 +1193,7 @@ internal fun parseConfig(root: YamlElement): Config {
think_tags = m.strOrNull("think_tags"),
backoff = m.durationOrNull("backoff"),
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
probe_interval = m.durationOrNull("probe_interval"),
)
}
@@ -10,6 +10,8 @@ import kotlinx.coroutines.sync.withLock
* Метрики (всё в памяти, сброс при рестарте):
* - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик
* 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) {
@@ -15,3 +15,12 @@ const val UPSTREAM_REQUEST_TIMEOUT_MS: Long = 5 * 60_000L
* Отключается явным `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
@@ -615,4 +615,32 @@ class ConfigLogicTest {
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,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")
}
}