feat: экспоненциальный backoff (providers[].backoff / upstreams[].backoff)
Повторяющиеся сбои апстрима (5xx/429/402 и сетевые ошибки) увязываются экспоненциальным откатом: первая ошибка — отдых 1s, далее 2s, 4s, … до потолка из конфига (ISO-8601: P1M, PT15M…); успешный запрос сбрасывает счётчик. Приоритет: upstreams[].backoff (свой счётчик на модель) над providers[].backoff (общий счётчик на все модели провайдера). - Backoff.kt: BackoffGuard (Mutex/@Volatile, kotlinx-datetime) + BackoffRegistry (скоупы провайдер/модель); - pickFreeUpstream пропускает апстримы в откате (лог «~Ns — пропускаю»); - все апстримы модели в откате → 503 «all upstreams cooling, retry in Ns» + заголовок Retry-After; - некорректный ISO-8601 в backoff — warn + игнор поля; - BackoffTest (8 тестов) + обновлённые тесты pickFreeUpstream; - CONFIG.md/README/TESTING.md. Проверено: jvmTest, compileKotlinLinuxX64, linuxX64Test.
This commit is contained in:
@@ -0,0 +1,111 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.datetime.Clock
|
||||
import kotlinx.datetime.Instant
|
||||
import kotlin.concurrent.Volatile
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
/**
|
||||
* Экспоненциальный бэкофф для одной области отдыха: первая ошибка — отдых
|
||||
* [base] (1 с), каждая следующая удваивает интервал вплоть до [cap].
|
||||
* Успешный запрос сбрасывает счётчик и отдых.
|
||||
*/
|
||||
class BackoffGuard(val cap: Duration, val base: Duration = 1.seconds) {
|
||||
private val lock = Mutex()
|
||||
@Volatile
|
||||
private var failures: Int = 0
|
||||
@Volatile
|
||||
private var interval: Duration = Duration.ZERO
|
||||
@Volatile
|
||||
private var coolingUntil: Instant = Instant.DISTANT_PAST
|
||||
|
||||
fun isCoolingAt(now: Instant): Boolean = now < coolingUntil
|
||||
|
||||
/** Сколько осталось до выхода из отдыха (0 — не в откате). */
|
||||
fun remainingAt(now: Instant): Duration =
|
||||
if (coolingUntil > now) coolingUntil - now else Duration.ZERO
|
||||
|
||||
/** Учёт ошибки; возвращает интервал, на который область уходит в отдых. */
|
||||
suspend fun recordFailure(now: Instant): Duration {
|
||||
lock.lock()
|
||||
try {
|
||||
failures++
|
||||
interval = if (failures == 1) base else interval * 2
|
||||
if (interval > cap) interval = cap
|
||||
coolingUntil = now + interval
|
||||
} finally {
|
||||
lock.unlock()
|
||||
}
|
||||
return interval
|
||||
}
|
||||
|
||||
/** Успех: сброс счётчика и отдыха. */
|
||||
suspend fun recordSuccess() {
|
||||
lock.lock()
|
||||
try {
|
||||
failures = 0
|
||||
interval = Duration.ZERO
|
||||
coolingUntil = Instant.DISTANT_PAST
|
||||
} finally {
|
||||
lock.unlock()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Реестр областей бэкофф-отдыха. Область апстрима:
|
||||
* - у модели (записи upstreams) задан `backoff` — свой счётчик и потолок (приоритет);
|
||||
* - иначе у провайдера задан `backoff` — общий для всех моделей провайдера счётчик;
|
||||
* - ни там, ни там — бэкофф для этого апстрима выключен.
|
||||
*/
|
||||
class BackoffRegistry(
|
||||
private val providers: Map<String, ProviderConf>,
|
||||
private val upstreams: Map<String, UpstreamConf>,
|
||||
) {
|
||||
private val upstreamGuards = mutableMapOf<String, BackoffGuard>()
|
||||
private val providerGuards = mutableMapOf<String, BackoffGuard>()
|
||||
private val lock = Mutex()
|
||||
|
||||
private fun guardForLocked(up: UpstreamConf): BackoffGuard? {
|
||||
up.backoff?.let { cap ->
|
||||
return upstreamGuards.getOrPut(up.id) { BackoffGuard(cap) }
|
||||
}
|
||||
providers[up.provider]?.backoff?.let { cap ->
|
||||
return providerGuards.getOrPut(up.provider) { BackoffGuard(cap) }
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
private suspend fun <T> locked(block: () -> T?): T? {
|
||||
lock.lock()
|
||||
return try {
|
||||
block()
|
||||
} finally {
|
||||
lock.unlock()
|
||||
}
|
||||
}
|
||||
|
||||
suspend fun isCooling(up: UpstreamConf): Boolean {
|
||||
val g = locked { guardForLocked(up) } ?: return false
|
||||
return g.isCoolingAt(Clock.System.now())
|
||||
}
|
||||
|
||||
/** Секунд до выхода апстрима из отдыха (0 — не в откате). */
|
||||
suspend fun remainingSeconds(up: UpstreamConf): Long {
|
||||
val g = locked { guardForLocked(up) } ?: return 0L
|
||||
return g.remainingAt(Clock.System.now()).inWholeSeconds
|
||||
}
|
||||
|
||||
/** Учёт ошибки; возвращает интервал отдыха, или null, если бэкофф для апстрима выключен. */
|
||||
suspend fun recordFailure(up: UpstreamConf): Duration? {
|
||||
val g = locked { guardForLocked(up) } ?: return null
|
||||
return g.recordFailure(Clock.System.now())
|
||||
}
|
||||
|
||||
/** Успех: сброс области апстрима. */
|
||||
suspend fun recordSuccess(up: UpstreamConf) {
|
||||
locked { guardForLocked(up) }?.recordSuccess()
|
||||
}
|
||||
}
|
||||
@@ -46,6 +46,7 @@ import net.mamoe.yamlkt.YamlList
|
||||
import net.mamoe.yamlkt.YamlLiteral
|
||||
import net.mamoe.yamlkt.YamlMap
|
||||
import kotlin.concurrent.Volatile
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
import kotlin.time.TimeSource
|
||||
|
||||
@@ -86,11 +87,18 @@ fun main() {
|
||||
val active = config.upstreams.associate { up ->
|
||||
up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider]))
|
||||
}
|
||||
val backoff = BackoffRegistry(providersById, upstreamsById)
|
||||
|
||||
log.info {
|
||||
"[llm-proxy] загружено: providers=${config.providers.size}, " +
|
||||
"upstreams=${config.upstreams.size}, models=${config.models.size} (config=$path)"
|
||||
}
|
||||
log.info {
|
||||
"[llm-proxy] backoff: upstreams=" +
|
||||
config.upstreams.filter { it.backoff != null }.map { up -> "${up.id}=${up.backoff?.toIsoString()}" }.joinToString(", ") +
|
||||
" providers=" +
|
||||
config.providers.filter { it.backoff != null }.map { p -> "${p.id}=${p.backoff?.toIsoString()}" }.joinToString(", ")
|
||||
}
|
||||
log.info {
|
||||
"[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}"
|
||||
}
|
||||
@@ -101,7 +109,7 @@ fun main() {
|
||||
val http = createHttpClient()
|
||||
val sessions = SessionRegistry()
|
||||
startServer(config.server.host, config.server.port) {
|
||||
proxyModule(config, providersById, upstreamsById, active, sessions, http)
|
||||
proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -114,10 +122,11 @@ fun Application.proxyModule(
|
||||
active: Map<String, UpstreamCounter>,
|
||||
sessions: SessionRegistry,
|
||||
http: HttpClient,
|
||||
backoff: BackoffRegistry,
|
||||
) {
|
||||
routing {
|
||||
post("/v1/chat/completions") {
|
||||
handleChat(call, config, providersById, upstreamsById, active, sessions, http)
|
||||
handleChat(call, config, providersById, upstreamsById, active, sessions, http, backoff)
|
||||
}
|
||||
get("/v1/models") {
|
||||
handleModels(call, config)
|
||||
@@ -133,6 +142,7 @@ private suspend fun handleChat(
|
||||
active: Map<String, UpstreamCounter>,
|
||||
sessions: SessionRegistry,
|
||||
http: HttpClient,
|
||||
backoff: BackoffRegistry,
|
||||
) {
|
||||
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
|
||||
val raw = call.receiveText()
|
||||
@@ -165,7 +175,7 @@ private suspend fun handleChat(
|
||||
var anyClaimed = false
|
||||
|
||||
while (true) {
|
||||
val up = pickFreeUpstream(pool, active, failed) ?: break
|
||||
val up = pickFreeUpstream(pool, active, failed, backoff, modelName) ?: break
|
||||
anyClaimed = true
|
||||
val start = TimeSource.Monotonic.markNow()
|
||||
try {
|
||||
@@ -230,7 +240,11 @@ private suspend fun handleChat(
|
||||
}.execute { resp ->
|
||||
upstreamStatus = resp.status.value
|
||||
if (upstreamStatus >= 500 || upstreamStatus == 429 || upstreamStatus == 402) {
|
||||
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" }
|
||||
val rest = backoff.recordFailure(up)
|
||||
log.warn {
|
||||
"[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" +
|
||||
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
|
||||
}
|
||||
failed.add(up.id)
|
||||
failover = true
|
||||
return@execute
|
||||
@@ -245,6 +259,7 @@ private suspend fun handleChat(
|
||||
return@execute
|
||||
}
|
||||
responded = true
|
||||
backoff.recordSuccess(up)
|
||||
if (clientWantsStream) {
|
||||
val ct = resp.headers["Content-Type"] ?: "text/event-stream"
|
||||
val status = HttpStatusCode.fromValue(upstreamStatus)
|
||||
@@ -294,7 +309,11 @@ private suspend fun handleChat(
|
||||
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" }
|
||||
val rest = backoff.recordFailure(up)
|
||||
log.error {
|
||||
"[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" +
|
||||
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
|
||||
}
|
||||
failed.add(up.id)
|
||||
continue
|
||||
} finally {
|
||||
@@ -305,7 +324,17 @@ private suspend fun handleChat(
|
||||
if (anyClaimed) {
|
||||
call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
|
||||
} else {
|
||||
call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
|
||||
val retryIn = pool.map { backoff.remainingSeconds(it) }.filter { it > 0 }.minOrNull()
|
||||
if (retryIn != null) {
|
||||
call.response.headers.append("Retry-After", retryIn.toString())
|
||||
call.respondText(
|
||||
errorJson("all upstreams cooling, retry in ${retryIn}s"),
|
||||
ContentType.Application.Json,
|
||||
HttpStatusCode.ServiceUnavailable,
|
||||
)
|
||||
} else {
|
||||
call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -489,15 +518,28 @@ class UpstreamCounter(private val limit: Int, initial: Int = 0) {
|
||||
|
||||
/**
|
||||
* Выбор апстрима для попытки: первый по порядку (приоритету) апстрим из `pool`,
|
||||
* у которого свободен слот и который ещё не в `excluded` (не упал ранее).
|
||||
* Сразу занимает слот (через [tryClaim]). Если свободных нет — возвращает null.
|
||||
* у которого свободен слот, который ещё не в `excluded` (не упал ранее) и
|
||||
* не в backoff-откате. Сразу занимает слот (через [tryClaim]).
|
||||
* Если свободных нет — возвращает null.
|
||||
*/
|
||||
internal fun pickFreeUpstream(
|
||||
internal suspend fun pickFreeUpstream(
|
||||
pool: List<UpstreamConf>,
|
||||
active: Map<String, UpstreamCounter>,
|
||||
excluded: Set<String>,
|
||||
): UpstreamConf? =
|
||||
pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) }
|
||||
backoff: BackoffRegistry,
|
||||
modelName: String = "?",
|
||||
): UpstreamConf? {
|
||||
for (up in pool) {
|
||||
if (up.id in excluded) continue
|
||||
val rest = backoff.remainingSeconds(up)
|
||||
if (rest > 0) {
|
||||
log.info { "[llm-proxy] model=$modelName upstream=${up.id} в backoff-откате (~${rest}c) — пропускаю" }
|
||||
continue
|
||||
}
|
||||
if (tryClaim(up, active)) return up
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
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)}" }
|
||||
@@ -903,6 +945,7 @@ data class ProviderConf(
|
||||
val think_tags: String? = null,
|
||||
val reasoning_field: String? = null,
|
||||
val reasoning_empty_ok: Boolean = false,
|
||||
val backoff: Duration? = null,
|
||||
)
|
||||
|
||||
data class UpstreamConf(
|
||||
@@ -912,6 +955,7 @@ data class UpstreamConf(
|
||||
val max_concurrency: Int? = null,
|
||||
val patch: JsonObject? = null,
|
||||
val think_tags: String? = null,
|
||||
val backoff: Duration? = null,
|
||||
)
|
||||
|
||||
data class ModelConf(
|
||||
@@ -969,6 +1013,7 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
think_tags = m.strOrNull("think_tags"),
|
||||
reasoning_field = m.strOrNull("reasoning_field"),
|
||||
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
|
||||
backoff = m.durationOrNull("backoff"),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -981,6 +1026,7 @@ internal fun parseConfig(root: YamlElement): Config {
|
||||
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
|
||||
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
|
||||
think_tags = m.strOrNull("think_tags"),
|
||||
backoff = m.durationOrNull("backoff"),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -1006,5 +1052,14 @@ 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. */
|
||||
internal fun Map<String, YamlElement>.durationOrNull(key: String): Duration? =
|
||||
strOrNull(key)?.let { raw ->
|
||||
runCatching { Duration.parse(raw) }.getOrElse {
|
||||
log.warn { "[llm-proxy] config: поле '$key' — некорректная ISO-8601-длительность '$raw', проигнорировано" }
|
||||
null
|
||||
}
|
||||
}
|
||||
|
||||
internal fun Map<String, YamlElement>.yamlMapOrNull(key: String): YamlMap? =
|
||||
this[key] as? YamlMap
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
package pw.binom.llmproxy
|
||||
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.datetime.Instant
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
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.milliseconds
|
||||
import kotlin.time.Duration.Companion.minutes
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
class BackoffTest {
|
||||
|
||||
private val t0 = Instant.fromEpochSeconds(1_700_000_000)
|
||||
|
||||
@Test
|
||||
fun intervalsDoubleUntilCap() = runTest {
|
||||
val g = BackoffGuard(cap = 30.seconds)
|
||||
val seq = (1..7).map { g.recordFailure(t0) }
|
||||
assertEquals(
|
||||
listOf(1.seconds, 2.seconds, 4.seconds, 8.seconds, 16.seconds, 30.seconds, 30.seconds),
|
||||
seq,
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun successResetsCounter() = runTest {
|
||||
val g = BackoffGuard(cap = 30.seconds)
|
||||
g.recordFailure(t0)
|
||||
g.recordFailure(t0)
|
||||
g.recordSuccess()
|
||||
assertEquals(1.seconds, g.recordFailure(t0))
|
||||
assertFalse(g.isCoolingAt(t0 + 1.days))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun coolingWindow() = runTest {
|
||||
val g = BackoffGuard(cap = 30.seconds)
|
||||
g.recordFailure(t0) // отдых 1s
|
||||
assertTrue(g.isCoolingAt(t0 + 500.milliseconds))
|
||||
assertFalse(g.isCoolingAt(t0 + 1.seconds))
|
||||
assertEquals(500.milliseconds, g.remainingAt(t0 + 500.milliseconds))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun iso8601CapParsing() {
|
||||
assertEquals(30.minutes, Duration.parse("PT30M"))
|
||||
assertEquals(1.days, Duration.parse("P1D"))
|
||||
assertEquals(90.seconds, Duration.parse("PT1M30S"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun providerScopeSharedBetweenItsModels() = runTest {
|
||||
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
|
||||
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
|
||||
val b = UpstreamConf(id = "b", provider = "p", model = "m2")
|
||||
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
|
||||
|
||||
assertEquals(1.seconds, reg.recordFailure(a))
|
||||
|
||||
// общая скоуп-переменная провайдера: сбой на a охладил и b
|
||||
assertTrue(reg.isCooling(a))
|
||||
assertTrue(reg.isCooling(b))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun upstreamBackoffOverridesProvider() = runTest {
|
||||
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
|
||||
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
|
||||
val b = UpstreamConf(id = "b", provider = "p", model = "m2", backoff = 30.seconds)
|
||||
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
|
||||
|
||||
// своё backoff у b: cap 30s (не 5m провайдера): 1s, 2s, 4s, 8s, 16s, 30s (потолок)
|
||||
val seq = (1..6).map { reg.recordFailure(b) }
|
||||
assertEquals(listOf(1.seconds, 2.seconds, 4.seconds, 8.seconds, 16.seconds, 30.seconds), seq)
|
||||
// скоупы независимы: сбои b не охлаждают a (провайдерский скоуп)
|
||||
assertFalse(reg.isCooling(a))
|
||||
assertTrue(reg.isCooling(b))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun providerFailureDoesNotCoolUpstreamWithOwnBackoff() = runTest {
|
||||
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
|
||||
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
|
||||
val b = UpstreamConf(id = "b", provider = "p", model = "m2", backoff = 30.minutes)
|
||||
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
|
||||
|
||||
// провайдерский скоуп заведён сбоем на a:
|
||||
reg.recordFailure(a)
|
||||
assertTrue(reg.isCooling(a))
|
||||
// b живёт своим (provider-scope на него не действует):
|
||||
assertFalse(reg.isCooling(b))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun noBackoffConfiguredMeansNoCooling() = runTest {
|
||||
val p = ProviderConf(id = "p", url = "https://p")
|
||||
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
|
||||
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a))
|
||||
|
||||
assertNull(reg.recordFailure(a))
|
||||
assertFalse(reg.isCooling(a))
|
||||
assertEquals(0L, reg.remainingSeconds(a))
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
import io.ktor.http.headersOf
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.json.JsonObject
|
||||
import kotlinx.serialization.json.jsonArray
|
||||
@@ -14,6 +15,8 @@ import net.mamoe.yamlkt.Yaml
|
||||
|
||||
class ConfigLogicTest {
|
||||
|
||||
private val noBackoff = BackoffRegistry(emptyMap(), emptyMap())
|
||||
|
||||
@Test
|
||||
fun mergeDeepMergesNestedObjectsAndReplacesScalars() {
|
||||
val base = Json.parseToJsonElement("""{"a":{"x":1,"y":2},"b":1}""").jsonObject
|
||||
@@ -222,58 +225,58 @@ class ConfigLogicTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() {
|
||||
fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() = runTest {
|
||||
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2))
|
||||
val pool = listOf(
|
||||
UpstreamConf("u1", "p", "m", 1, null),
|
||||
UpstreamConf("u2", "p", "m", 2, null),
|
||||
)
|
||||
val up = pickFreeUpstream(pool, active, emptySet())
|
||||
val up = pickFreeUpstream(pool, active, emptySet(), noBackoff)
|
||||
assertEquals("u1", up?.id)
|
||||
// слот реально занят
|
||||
assertEquals(1, active.getValue("u1").current)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun pickFreeUpstreamSkipsExcluded() {
|
||||
fun pickFreeUpstreamSkipsExcluded() = runTest {
|
||||
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2))
|
||||
val pool = listOf(
|
||||
UpstreamConf("u1", "p", "m", 1, null),
|
||||
UpstreamConf("u2", "p", "m", 2, null),
|
||||
)
|
||||
val up = pickFreeUpstream(pool, active, setOf("u1"))
|
||||
val up = pickFreeUpstream(pool, active, setOf("u1"), noBackoff)
|
||||
assertEquals("u2", up?.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun pickFreeUpstreamReturnsNullWhenAllBusy() {
|
||||
fun pickFreeUpstreamReturnsNullWhenAllBusy() = runTest {
|
||||
val active = mapOf("u1" to UpstreamCounter(1, 1)) // уже на лимите 1
|
||||
val pool = listOf(UpstreamConf("u1", "p", "m", 1, null))
|
||||
assertEquals(null, pickFreeUpstream(pool, active, emptySet()))
|
||||
assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun pickFreeUpstreamImplementsFailoverOrder() {
|
||||
fun pickFreeUpstreamImplementsFailoverOrder() = runTest {
|
||||
// dead исключён (упал ранее) — выбирается следующий живой u1
|
||||
val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1))
|
||||
val pool = listOf(
|
||||
UpstreamConf("dead", "p", "m", 1, null),
|
||||
UpstreamConf("u1", "p", "m", 1, null),
|
||||
)
|
||||
val up = pickFreeUpstream(pool, active, setOf("dead"))
|
||||
val up = pickFreeUpstream(pool, active, setOf("dead"), noBackoff)
|
||||
assertEquals("u1", up?.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() {
|
||||
fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() = runTest {
|
||||
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1))
|
||||
val pool = listOf(
|
||||
UpstreamConf("u1", "p", "m", 1, null),
|
||||
UpstreamConf("u2", "p", "m", 1, null),
|
||||
)
|
||||
assertEquals("u1", pickFreeUpstream(pool, active, emptySet())?.id)
|
||||
assertEquals("u2", pickFreeUpstream(pool, active, emptySet())?.id)
|
||||
assertEquals(null, pickFreeUpstream(pool, active, emptySet())?.id)
|
||||
assertEquals("u1", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
|
||||
assertEquals("u2", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
|
||||
assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user