diff --git a/CONFIG.md b/CONFIG.md index c31e49b..637aea9 100644 --- a/CONFIG.md +++ b/CONFIG.md @@ -322,6 +322,69 @@ 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= upstream= TIMEOUT: первый байт не пришёл за 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`. + ### Prometheus-метрики (`/metrics`) Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics` @@ -331,7 +394,7 @@ 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_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 +406,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 +492,10 @@ 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) ) -@Serializable data class UpstreamConf( val id: String, // наш внутренний id апстрима val provider: String, // ссылка на ProviderConf.id @@ -441,6 +504,7 @@ data class UpstreamConf( val patch: JsonObject? = null, // к запросам этой апстрим-модели val think_tags: String? = null, // переопределение think-режима модели val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским) + val first_byte_timeout: Duration? = null, // окно первого байта стрима (приоритет над провайдерским) ) @Serializable @@ -605,6 +669,10 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet() (`(backoff: отдых PT…)`); при пропуске апстрима в откате — `upstream= в backoff-откате (~Ns) — пропускаю`; при старте — список настроенных backoff (`backoff: upstreams=… providers=…`). +- **Таймаут первого байта стрима**: при срабатывании watchdog'а — + `upstream= TIMEOUT: первый байт не пришёл за ms → фейловер`; + при старте — список настроенных `first_byte_timeout` + (`first_byte_timeout: upstreams=… providers=… default=PT30S`). - **Все в откате** — `503 all upstreams cooling, retry in Ns` для `model=<витрина>` с заголовком `Retry-After: N`. diff --git a/README.md b/README.md index af23228..2a3f326 100644 --- a/README.md +++ b/README.md @@ -30,6 +30,12 @@ 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`. + Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по `model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) — подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt index 6bd865d..466401b 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt @@ -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() 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) { 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: \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() 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 { + 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 { + 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"), ) } diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt index ee410e6..fb27f57 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt @@ -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); diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt index 9c104b5..9c330f9 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt @@ -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 diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt index b5a9443..71979bb 100644 --- a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt +++ b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt @@ -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,32 @@ 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) + } } diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/FirstByteTimeoutTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/FirstByteTimeoutTest.kt new file mode 100644 index 0000000..d4b2514 --- /dev/null +++ b/src/commonTest/kotlin/pw/binom/llmproxy/FirstByteTimeoutTest.kt @@ -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() + + 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() + + 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() + + 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() + + 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)) + } +} \ No newline at end of file