fix(client): отключаем request/connect/socket-таймауты для SSE-стримов
ci / JVM build + tests (push) Failing after 1m23s

Дефолтный CIOEngineConfig.requestTimeout = 15 с убивал SSE-стрим при
простое, потому что движок CIO не считает запрос SSE-шным (мы читаем
bodyAsChannel() руками, без SSEClientContent). На TUI это проявлялось как
'стрим отвалился через 15 с' — события молча переставали приходить.

Два уровня фикса:

1. Per-request: HttpRequestBuilder.noSseReadTimeout() ставит capability
   HttpTimeoutCapability со всеми таймаутами = INFINITE_TIMEOUT_MS.
   В ConversationClient.events() и AgentClient.events() вызывается перед
   каждым SSE-стримом. Плагин HttpTimeout (если установлен) читает эту
   capability через ?: и не перезаписывает её.

2. Default client: defaultAgentikHttpClient() ставит
   engine { requestTimeout = 0 } — defense-in-depth на случай, если
   кто-то соберёт свой HttpClient без capability.

Тесты:
- SseTimeoutTest запускает встроенный Ktor CIO-сервер, держит stream 17 с.
- 'with noSseReadTimeout' — stream живёт до 'done' (тест проходит ~17 с).
- 'without noSseReadTimeout' — клиент падает на ~15 с с
  HttpRequestTimeoutException (контр-тест, доказывает что баг был).
This commit is contained in:
2026-09-17 14:39:01 +03:00
parent db3c49099c
commit 9d826a4e81
6 changed files with 207 additions and 2 deletions
@@ -66,7 +66,9 @@ internal class AgentClient(
}
override fun events(after: Instant): Flow<AgentEvent> = flow {
val response = httpClient.get("$agentUrl/events?after=$after")
val response = httpClient.get("$agentUrl/events?after=$after") {
noSseReadTimeout()
}
check(response.status == HttpStatusCode.OK) {
"events: server returned ${response.status}"
}
@@ -37,7 +37,20 @@ fun AgentikAgent(
* Дефолтный [HttpClient] для общения с `agentikAgent`: CIO-движок и
* kotlinx-serialization с тем же wire-форматом, что на сервере. SSE-парсер
* (см. [readSse]) живёт в общем коде и плагина не требует.
*
* Engine config: `requestTimeout = 0` отключает встроенный request-таймаут CIO
* (по умолчанию 15 с — наш кастомный SSE-ридер не маркирует запрос как
* `SSEClientContent`, и `getRequestTimeout` для движка вернул бы 15000 мс).
* Сами SSE-запросы в `ConversationClient.events`/`AgentClient.events` и так
* уже выставляют `HttpTimeoutCapability` = INFINITE (см. [noSseReadTimeout]),
* но `engine { requestTimeout = 0 }` — defense-in-depth: если кто-то
* использует свой HttpClient через `AgentikAgent(..., httpClient = ...)` без
* capability, фоновый стрим на нашем URL тоже не убьётся по 15-секундному
* таймеру.
*
* Connect/socket таймауты оставлены дефолтными (5 с / бесконечность).
*/
fun defaultAgentikHttpClient(): HttpClient = HttpClient(CIO) {
engine { requestTimeout = 0 }
install(ContentNegotiation) { json(agentikJson) }
}
@@ -72,7 +72,9 @@ internal class ConversationClient(
}
override fun events(after: Instant): Flow<Event> = flow {
val response = httpClient.get("$convUrl/events?after=$after")
val response = httpClient.get("$convUrl/events?after=$after") {
noSseReadTimeout()
}
check(response.status == HttpStatusCode.OK) {
"events: server returned ${response.status}"
}
@@ -0,0 +1,31 @@
package pw.binom.agentik.client
import io.ktor.client.plugins.HttpTimeoutConfig
import io.ktor.client.plugins.HttpTimeoutCapability
import io.ktor.client.request.HttpRequestBuilder
/**
* Отключает request/connect/socket-таймауты для конкретного запроса через
* [HttpTimeoutCapability] со всеми таймаутами = [HttpTimeoutConfig.INFINITE_TIMEOUT_MS].
*
* Зачем: наш SSE-ридер ([readSse]) читает `bodyAsChannel()` руками и не
* использует плагин `SSE`, поэтому движок не считает запрос SSE-шным
* (`HttpRequestBuilder.supportsRequestTimeout` проверяет
* `body is SSEClientContent`, а у нас тело — обычный GET без тела).
* Без capability встроенный `CIOEngineConfig.requestTimeout` (по умолчанию
* **15000 мс**) молча убивает долгий idle-стрим через 15 секунд.
*
* Конфиг создаётся заново на каждый вызов — плагин `HttpTimeout` при
* установленном capability мутирует его поля через `?:`, так что шаренный
* инстанс мог бы утечь между запросами.
*/
internal fun HttpRequestBuilder.noSseReadTimeout() {
setCapability(
HttpTimeoutCapability,
HttpTimeoutConfig(
requestTimeoutMillis = HttpTimeoutConfig.INFINITE_TIMEOUT_MS,
connectTimeoutMillis = HttpTimeoutConfig.INFINITE_TIMEOUT_MS,
socketTimeoutMillis = HttpTimeoutConfig.INFINITE_TIMEOUT_MS,
),
)
}
@@ -0,0 +1,148 @@
package pw.binom.agentik.client
import io.ktor.client.HttpClient
import io.ktor.client.engine.cio.CIO
import io.ktor.client.plugins.HttpRequestTimeoutException
import io.ktor.client.request.header
import io.ktor.client.request.prepareGet
import io.ktor.client.statement.bodyAsChannel
import io.ktor.server.application.call
import io.ktor.server.engine.embeddedServer
import io.ktor.server.response.respondBytesWriter
import io.ktor.server.routing.get
import io.ktor.server.routing.routing
import io.ktor.http.ContentType
import io.ktor.utils.io.writeStringUtf8
import io.ktor.utils.io.readUTF8Line
import kotlinx.coroutines.delay
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeout
import java.net.ServerSocket
import kotlin.test.Test
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
import kotlin.test.fail
/**
* Репродукция бага Ktor CIO: дефолтный [io.ktor.client.engine.cio.CIOEngineConfig.requestTimeout]
* = 15 с убивает SSE read. Наш fix — [noSseReadTimeout] ставит capability
* [io.ktor.client.plugins.HttpTimeoutCapability] со всеми таймаутами = INFINITE
* перед каждым read-стримом.
*
* Тест запускает встроенный Ktor CIO-сервер на свободном порту. Сервер шлёт
* "hello", ждёт 20 с (дольше дефолтного requestTimeout = 15 с), затем шлёт
* "done". Без capability клиент отвалился бы на ~15 с; с capability — второе
* сообщение доходит.
*
* Читаем строки пока не найдём "data: done" или пока не сработает
* [withTimeout] (18 с — запас над server delay 20 с).
*/
class SseTimeoutTest {
private fun freePort(): Int = ServerSocket(0).use { it.localPort }
@Test
fun `sse read survives past default cio timeout with noSseReadTimeout`(): Unit = runBlocking {
val port = freePort()
val server = embeddedServer(io.ktor.server.cio.CIO, port = port) {
routing {
get("/sse") {
call.respondBytesWriter(contentType = ContentType.Text.EventStream) {
writeStringUtf8("data: hello\n\n")
flush()
// 17 с — чуть больше дефолтного CIO requestTimeout = 15 с.
// Если capability сломана, клиент упадёт на 15 с и не получит "done".
delay(17_000)
writeStringUtf8("data: done\n\n")
}
}
}
}.start(wait = false)
try {
val client = HttpClient(CIO)
val received = mutableListOf<String>()
client.prepareGet("http://127.0.0.1:$port/sse") {
header("Accept", "text/event-stream")
noSseReadTimeout()
}.execute { resp ->
val ch = resp.bodyAsChannel()
// 19 с запас: ждём, пока сервер пошлёт "done" после 17 с.
// Если capability сломана, клиент упадёт на 15 с и мы словим исключение.
val deadline = 19_000L
val start = System.currentTimeMillis()
while (System.currentTimeMillis() - start < deadline) {
val line = withTimeout<String?>(deadline) { ch.readUTF8Line() } ?: break
if (line.startsWith("data: ")) {
received.add(line)
}
if (line == "data: done") break
}
}
assertTrue(received.contains("data: hello"), "должно получить hello: $received")
assertTrue(
received.contains("data: done"),
"должно получить done (SSE read не должен падать на 15 с): $received",
)
assertFalse(
received.any { it == "<timeout>" },
"SSE read упал в timeout (capability не сработал): $received",
)
} finally {
server.stop(100, 200)
}
}
/**
* Контр-тест: убеждаемся что БЕЗ [noSseReadTimeout] дефолтный
* CIO requestTimeout = 15 с действительно убивает SSE-стрим.
* Сервер держит stream 17 с; если клиент не выставил capability —
* мы должны получить [HttpRequestTimeoutException] на ~15 с, не
* дожидаясь "done".
*/
@Test
fun `without noSseReadTimeout default cio requestTimeout kills the stream`(): Unit = runBlocking {
val port = freePort()
val server = embeddedServer(io.ktor.server.cio.CIO, port = port) {
routing {
get("/sse") {
call.respondBytesWriter(contentType = ContentType.Text.EventStream) {
writeStringUtf8("data: hello\n\n")
flush()
delay(17_000)
writeStringUtf8("data: done\n\n")
}
}
}
}.start(wait = false)
try {
val client = HttpClient(CIO)
val start = System.currentTimeMillis()
try {
client.prepareGet("http://127.0.0.1:$port/sse") {
header("Accept", "text/event-stream")
// НАМЕРЕННО без noSseReadTimeout.
}.execute { resp ->
val ch = resp.bodyAsChannel()
// Читаем строки, пока не придёт "data: done" — без capability
// клиент упадёт на ~15 с до того, как сервер пошлёт done.
while (true) {
val line = ch.readUTF8Line() ?: break
if (line == "data: done") break
}
}
fail("без capability клиент должен словить HttpRequestTimeoutException")
} catch (e: HttpRequestTimeoutException) {
val elapsed = System.currentTimeMillis() - start
assertTrue(
elapsed in 14_000..17_000,
"timeout должен сработать в районе 15 с (default), elapsed=$elapsed",
)
}
} finally {
server.stop(100, 200)
}
}
}