From de95013164c5dc3adaaa05950d751bc7e19aa321 Mon Sep 17 00:00:00 2001 From: Liuxinyu176 <1041316040@qq.com> Date: Thu, 8 Oct 2026 23:52:12 +0800 Subject: [PATCH] =?UTF-8?q?feat(proxy):=20stream=3Dtrue=20=E5=AE=9E?= =?UTF-8?q?=E6=97=B6=E9=80=8F=E4=BC=A0=E4=B8=8A=E6=B8=B8=20SSE=20=E5=AD=97?= =?UTF-8?q?=E8=8A=82=E6=B5=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Ktor respondOutputStream 逐块转发上游 body,不再整包缓冲 - 流式/非流式按客户端 body 的 stream 字段分流 - openStreamingChat 入口正式接线 --- .../token/data/proxy/KtorLocalProxyServer.kt | 63 +++++++++++++++---- 1 file changed, 50 insertions(+), 13 deletions(-) diff --git a/app/src/main/java/com/rainy/token/data/proxy/KtorLocalProxyServer.kt b/app/src/main/java/com/rainy/token/data/proxy/KtorLocalProxyServer.kt index aacdc13..48be89a 100644 --- a/app/src/main/java/com/rainy/token/data/proxy/KtorLocalProxyServer.kt +++ b/app/src/main/java/com/rainy/token/data/proxy/KtorLocalProxyServer.kt @@ -14,6 +14,7 @@ import io.ktor.server.plugins.contentnegotiation.ContentNegotiation import io.ktor.server.request.receiveText import io.ktor.server.response.respond import io.ktor.server.response.respondBytes +import io.ktor.server.response.respondOutputStream import io.ktor.server.routing.get import io.ktor.server.routing.post import io.ktor.server.routing.routing @@ -129,22 +130,58 @@ class KtorLocalProxyServer @Inject constructor( val pooled = pool.next(route.kind, route.region, conversationId) val accountId = pooled?.accountId try { - val result: ProxyUpstreamResponse? = when (route.kind) { - ProviderKind.WORKBUDDY_CN, ProviderKind.WORKBUDDY_INTL -> - workBuddy.forwardChat(body, accountId, route.region) + if (extractStream(body)) { + val stream: ProxyUpstreamStream? = when (route.kind) { + ProviderKind.WORKBUDDY_CN, ProviderKind.WORKBUDDY_INTL -> + workBuddy.openStreamingChat(body, accountId, route.region) - ProviderKind.TRAE_CN, ProviderKind.TRAE_INTL -> - trae.forwardChat(body, accountId, route.region) + ProviderKind.TRAE_CN, ProviderKind.TRAE_INTL -> + trae.openStreamingChat(body, accountId, route.region) - else -> sub2Api.forwardChat(body, accountId) - } - if (result == null) { - call.respond( - HttpStatusCode.BadRequest, - errorBody("${route.kind.displayName} 未配置或未登录,请先在设置中配置") - ) + else -> sub2Api.openStreamingChat(body, accountId) + } + if (stream == null) { + call.respond( + HttpStatusCode.BadRequest, + errorBody("${route.kind.displayName} 未配置或未登录,请先在设置中配置") + ) + } else { + call.respondOutputStream( + contentType = contentTypeOf(stream.contentType), + status = HttpStatusCode(stream.status, "") + ) { + try { + val buffer = ByteArray(8192) + val input = stream.input + while (true) { + val read = input.read(buffer) + if (read < 0) break + write(buffer, 0, read) + flush() + } + } finally { + stream.close() + } + } + } } else { - call.respondBytes(result.body, contentTypeOf(result.contentType), HttpStatusCode(result.status, "")) + val result: ProxyUpstreamResponse? = when (route.kind) { + ProviderKind.WORKBUDDY_CN, ProviderKind.WORKBUDDY_INTL -> + workBuddy.forwardChat(body, accountId, route.region) + + ProviderKind.TRAE_CN, ProviderKind.TRAE_INTL -> + trae.forwardChat(body, accountId, route.region) + + else -> sub2Api.forwardChat(body, accountId) + } + if (result == null) { + call.respond( + HttpStatusCode.BadRequest, + errorBody("${route.kind.displayName} 未配置或未登录,请先在设置中配置") + ) + } else { + call.respondBytes(result.body, contentTypeOf(result.contentType), HttpStatusCode(result.status, "")) + } } } catch (e: IOException) { val detail = e.message ?: "未知错误"