feat: watchdog на первый байт стрима (first_byte_timeout) — фейловер + backoff при молчании апстрима

This commit is contained in:
2026-09-15 00:37:51 +03:00
parent 293964d8f1
commit 249b83eb51
7 changed files with 465 additions and 11 deletions
+126 -6
View File
@@ -28,6 +28,11 @@ 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.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 +108,13 @@ 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] server: bind host=${config.server.host} port=${config.server.port}"
}
@@ -279,13 +291,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 +335,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 +515,27 @@ 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
}
/**
* Сигнализирует, что апстрим не прислал ни одного байта тела стрим-ответа за
* [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 +870,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 +910,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 +986,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 +1082,7 @@ data class ProviderConf(
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
)
data class UpstreamConf(
@@ -976,6 +1093,7 @@ data class UpstreamConf(
val patch: JsonObject? = null,
val think_tags: String? = null,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
)
data class ModelConf(
@@ -1034,6 +1152,7 @@ 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"),
)
}
@@ -1047,6 +1166,7 @@ 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"),
)
}
@@ -9,7 +9,7 @@ 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_upstream_inflight{upstream,provider}` — сейчас в работе (слоты);
* - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев
* (счётчик backoff);
@@ -1,4 +1,17 @@
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