diff --git a/client/build.gradle.kts b/client/build.gradle.kts index a2c3d5c..db83f7b 100644 --- a/client/build.gradle.kts +++ b/client/build.gradle.kts @@ -22,4 +22,13 @@ dependencies { implementation(libs.kotlinx.coroutines.core) implementation(libs.kotlinx.serialization.core) 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) } diff --git a/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt b/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt index 15423a7..ac3ee35 100644 --- a/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt +++ b/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -66,7 +66,9 @@ internal class AgentClient( } override fun events(after: Instant): Flow = 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}" } diff --git a/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt b/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt index 6737836..7f92e37 100644 --- a/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt +++ b/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt @@ -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) } } diff --git a/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt b/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt index cd7a783..390fcb9 100644 --- a/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt +++ b/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt @@ -72,7 +72,9 @@ internal class ConversationClient( } override fun events(after: Instant): Flow = 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}" } diff --git a/client/src/main/kotlin/pw/binom/agentik/client/SseTimeout.kt b/client/src/main/kotlin/pw/binom/agentik/client/SseTimeout.kt new file mode 100644 index 0000000..d15e7d5 --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/SseTimeout.kt @@ -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, + ), + ) +} diff --git a/client/src/test/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt b/client/src/test/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt new file mode 100644 index 0000000..67b9b6c --- /dev/null +++ b/client/src/test/kotlin/pw/binom/agentik/client/SseTimeoutTest.kt @@ -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() + 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(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 == "" }, + "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) + } + } +}