fix(client): отключаем request/connect/socket-таймауты для SSE-стримов
ci / JVM build + tests (push) Failing after 1m23s
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:
@@ -22,4 +22,13 @@ dependencies {
|
|||||||
implementation(libs.kotlinx.coroutines.core)
|
implementation(libs.kotlinx.coroutines.core)
|
||||||
implementation(libs.kotlinx.serialization.core)
|
implementation(libs.kotlinx.serialization.core)
|
||||||
implementation(libs.kotlinx.serialization.json)
|
implementation(libs.kotlinx.serialization.json)
|
||||||
|
|
||||||
|
testImplementation(libs.kotlin.test)
|
||||||
|
testImplementation(libs.kotlinx.coroutines.test)
|
||||||
|
testImplementation("junit:junit:4.13.2")
|
||||||
|
testImplementation(libs.ktor.server.core)
|
||||||
|
testImplementation(libs.ktor.server.cio)
|
||||||
|
testImplementation(libs.ktor.server.test.host)
|
||||||
|
testImplementation(libs.ktor.server.sse)
|
||||||
|
testImplementation(libs.ktor.client.content.negotiation)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -66,7 +66,9 @@ internal class AgentClient(
|
|||||||
}
|
}
|
||||||
|
|
||||||
override fun events(after: Instant): Flow<AgentEvent> = flow {
|
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) {
|
check(response.status == HttpStatusCode.OK) {
|
||||||
"events: server returned ${response.status}"
|
"events: server returned ${response.status}"
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -37,7 +37,20 @@ fun AgentikAgent(
|
|||||||
* Дефолтный [HttpClient] для общения с `agentikAgent`: CIO-движок и
|
* Дефолтный [HttpClient] для общения с `agentikAgent`: CIO-движок и
|
||||||
* kotlinx-serialization с тем же wire-форматом, что на сервере. SSE-парсер
|
* kotlinx-serialization с тем же wire-форматом, что на сервере. SSE-парсер
|
||||||
* (см. [readSse]) живёт в общем коде и плагина не требует.
|
* (см. [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) {
|
fun defaultAgentikHttpClient(): HttpClient = HttpClient(CIO) {
|
||||||
|
engine { requestTimeout = 0 }
|
||||||
install(ContentNegotiation) { json(agentikJson) }
|
install(ContentNegotiation) { json(agentikJson) }
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -72,7 +72,9 @@ internal class ConversationClient(
|
|||||||
}
|
}
|
||||||
|
|
||||||
override fun events(after: Instant): Flow<Event> = flow {
|
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) {
|
check(response.status == HttpStatusCode.OK) {
|
||||||
"events: server returned ${response.status}"
|
"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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user