From 85c6cf428331eb62c70ef03e655332afc373ca8b Mon Sep 17 00:00:00 2001 From: Liuxinyu176 <1041316040@qq.com> Date: Thu, 8 Oct 2026 22:58:33 +0800 Subject: [PATCH] =?UTF-8?q?feat(proxy):=20M1a=20=E6=9C=AC=E5=9C=B0=20HTTP?= =?UTF-8?q?=20=E6=9C=8D=E5=8A=A1=20+=20Sub2API=20=E9=80=8F=E4=BC=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 引入 Ktor 3.0.3(server-core / cio / content-negotiation / kotlinx-json) - Sub2ApiChatProxy:把 /v1/chat/completions 与 /v1/models 透传到用户配置的 Sub2API 实例 - KtorLocalProxyServer:绑定 127.0.0.1:8787,提供 /health、/v1/models、/v1/chat/completions - NetworkModule 提供 Hilt 注入 --- app/build.gradle.kts | 6 + .../token/data/proxy/KtorLocalProxyServer.kt | 109 ++++++++++++++++++ .../token/data/proxy/Sub2ApiChatProxy.kt | 103 +++++++++++++++++ .../java/com/rainy/token/di/NetworkModule.kt | 18 +++ gradle/libs.versions.toml | 7 ++ 5 files changed, 243 insertions(+) create mode 100644 app/src/main/java/com/rainy/token/data/proxy/KtorLocalProxyServer.kt create mode 100644 app/src/main/java/com/rainy/token/data/proxy/Sub2ApiChatProxy.kt diff --git a/app/build.gradle.kts b/app/build.gradle.kts index 077c966..1726179 100644 --- a/app/build.gradle.kts +++ b/app/build.gradle.kts @@ -214,6 +214,12 @@ dependencies { implementation(libs.kotlinx.serialization.json) implementation(libs.retrofit.kotlinx.serialization.converter) + // Ktor 本地反代 HTTP 服务 + implementation(libs.ktor.server.core) + implementation(libs.ktor.server.cio) + implementation(libs.ktor.server.content.negotiation) + implementation(libs.ktor.serialization.kotlinx.json) + // DataStore implementation(libs.androidx.datastore.preferences) 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 new file mode 100644 index 0000000..fe90253 --- /dev/null +++ b/app/src/main/java/com/rainy/token/data/proxy/KtorLocalProxyServer.kt @@ -0,0 +1,109 @@ +package com.rainy.token.data.proxy + +import io.ktor.http.ContentType +import io.ktor.http.HttpStatusCode +import io.ktor.serialization.kotlinx.json.json +import io.ktor.server.application.Application +import io.ktor.server.application.install +import io.ktor.server.cio.CIO +import io.ktor.server.engine.EmbeddedServer +import io.ktor.server.engine.embeddedServer +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.routing.get +import io.ktor.server.routing.post +import io.ktor.server.routing.routing +import javax.inject.Inject +import javax.inject.Singleton +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow + +/** + * 基于 Ktor CIO 的本地 HTTP 反代服务。 + * + * M1a 能力: + * - GET /health + * - GET /v1/models → Sub2API 透传 + * - POST /v1/chat/completions → Sub2API 透传 + * + * 只绑定 127.0.0.1,默认端口 8787。 + */ +@Singleton +class KtorLocalProxyServer @Inject constructor( + private val sub2ApiChatProxy: Sub2ApiChatProxy, +) : LocalProxyServer { + + private val _isRunning = MutableStateFlow(false) + override val isRunning: StateFlow = _isRunning.asStateFlow() + + private var server: EmbeddedServer<*, *>? = null + + override fun start(config: ProxyServerConfig): Result = try { + if (_isRunning.value) return Result.success(Unit) + val engine = embeddedServer(CIO, host = "127.0.0.1", port = config.port) { + proxyModule(config, sub2ApiChatProxy) + } + engine.start(wait = false) + server = engine + _isRunning.value = true + Result.success(Unit) + } catch (e: Throwable) { + Result.failure(e) + } + + override fun stop() { + runCatching { server?.stop(gracePeriodMillis = 500, timeoutMillis = 2000) } + server = null + _isRunning.value = false + } + + private fun Application.proxyModule( + config: ProxyServerConfig, + sub2Api: Sub2ApiChatProxy, + ) { + install(ContentNegotiation) { + json() + } + routing { + get("/health") { + call.respond(mapOf("status" to "ok")) + } + get("/v1/models") { + val result = sub2Api.forwardModels() + if (result == null) { + call.respondBytes( + errorBody("Sub2API 未配置或未登录,请在设置中填写 API Key"), + ContentType.Application.Json, + HttpStatusCode.BadRequest + ) + } else { + call.respondBytes(result.body, contentTypeOf(result.contentType), HttpStatusCode(result.status, "")) + } + } + post("/v1/chat/completions") { + val body = call.receiveText() + val result = sub2Api.forwardChat(body) + if (result == null) { + call.respondBytes( + errorBody("Sub2API 未配置或未登录,请在设置中填写 API Key"), + ContentType.Application.Json, + HttpStatusCode.BadRequest + ) + } else { + call.respondBytes(result.body, contentTypeOf(result.contentType), HttpStatusCode(result.status, "")) + } + } + } + } + + private fun errorBody(message: String): ByteArray { + val safe = message.replace("\"", "'") + return "{\"error\":{\"message\":\"$safe\",\"type\":\"invalid_request_error\"}}".toByteArray() + } + + private fun contentTypeOf(raw: String): ContentType = + runCatching { ContentType.parse(raw) }.getOrDefault(ContentType.Application.Json) +} diff --git a/app/src/main/java/com/rainy/token/data/proxy/Sub2ApiChatProxy.kt b/app/src/main/java/com/rainy/token/data/proxy/Sub2ApiChatProxy.kt new file mode 100644 index 0000000..2048d13 --- /dev/null +++ b/app/src/main/java/com/rainy/token/data/proxy/Sub2ApiChatProxy.kt @@ -0,0 +1,103 @@ +package com.rainy.token.data.proxy + +import com.rainy.token.data.repository.CredentialRepository +import com.rainy.token.domain.model.Credential +import com.rainy.token.domain.service.ServiceType +import javax.inject.Inject +import javax.inject.Singleton +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.withContext +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.RequestBody.Companion.toRequestBody + +/** + * Sub2API Chat 透传代理。 + * + * Sub2API 实例本身暴露 OpenAI 兼容接口,因此这里不做协议转换: + * 取出用户保存的 Sub2ApiCredential(优先 sk- API Key,其次面板 authToken), + * 把客户端请求原样转发到 {base}/v1/chat/completions, + * 并把上游响应(含 SSE 流式内容)原样返回给本地客户端。 + * + * M1a 阶段:先做整包透传(流式也先缓冲),后续由 StreamNormalizer 升级为逐块转发。 + */ +@Singleton +class Sub2ApiChatProxy @Inject constructor( + private val okHttpClient: OkHttpClient, + private val credentialRepository: CredentialRepository, +) { + + /** 转发 POST /v1/chat/completions。 */ + suspend fun forwardChat(requestBody: String, accountId: String? = null): ProxyUpstreamResponse? { + val base = resolveBase(accountId) ?: return null + return forward(base = base, path = "/v1/chat/completions", requestBody = requestBody, accountId = accountId) + } + + /** 转发 GET /v1/models。 */ + suspend fun forwardModels(accountId: String? = null): ProxyUpstreamResponse? { + val base = resolveBase(accountId) ?: return null + return forward(base = base, path = "/v1/models", requestBody = null, accountId = accountId) + } + + private suspend fun resolveBase(accountId: String?): String? { + val credential = credentialRepository.get(ServiceType.SUB2API, accountId) ?: return null + if (credential !is Credential.Sub2ApiCredential) return null + if (resolveAuth(credential) == null) return null + return normalizeBase(credential.baseUrl) + } + + private suspend fun forward( + base: String, + path: String, + requestBody: String?, + accountId: String?, + ): ProxyUpstreamResponse? = withContext(Dispatchers.IO) { + val credential = credentialRepository.get(ServiceType.SUB2API, accountId) + ?: return@withContext null + if (credential !is Credential.Sub2ApiCredential) return@withContext null + val auth = resolveAuth(credential) ?: return@withContext null + + val builder = Request.Builder() + .url(base + path) + .addHeader("Authorization", auth) + val body = requestBody?.takeIf { it.isNotBlank() } + if (body != null) { + builder + .addHeader("Content-Type", "application/json") + .post(body.toRequestBody("application/json".toMediaType())) + } + + val response = try { + okHttpClient.newCall(builder.build()).execute() + } catch (_: Throwable) { + return@withContext null + } + val bytes = try { response.body?.bytes() ?: ByteArray(0) } catch (_: Throwable) { ByteArray(0) } + val contentType = response.header("Content-Type") ?: "application/json" + val status = response.code + response.close() + ProxyUpstreamResponse(status, contentType, bytes) + } + + private fun resolveAuth(credential: Credential.Sub2ApiCredential): String? { + credential.apiKey?.trim()?.takeIf { it.isNotBlank() }?.let { return "Bearer $it" } + credential.authToken?.trim()?.takeIf { it.isNotBlank() }?.let { return "Bearer $it" } + return null + } + + private fun normalizeBase(raw: String): String? { + var s = raw.trim() + while (s.endsWith("/")) s = s.dropLast(1) + return s.takeIf { it.isNotBlank() } + } +} + +/** + * 上游 HTTP 响应(透传用)。 + */ +data class ProxyUpstreamResponse( + val status: Int, + val contentType: String, + val body: ByteArray, +) diff --git a/app/src/main/java/com/rainy/token/di/NetworkModule.kt b/app/src/main/java/com/rainy/token/di/NetworkModule.kt index c89cf21..9af3644 100644 --- a/app/src/main/java/com/rainy/token/di/NetworkModule.kt +++ b/app/src/main/java/com/rainy/token/di/NetworkModule.kt @@ -20,6 +20,9 @@ import com.rainy.token.data.repository.OllamaRepository import com.rainy.token.data.repository.Sub2ApiRepository import com.rainy.token.data.repository.TraeRepository import com.rainy.token.data.repository.UpdateRepository +import com.rainy.token.data.proxy.KtorLocalProxyServer +import com.rainy.token.data.proxy.LocalProxyServer +import com.rainy.token.data.proxy.Sub2ApiChatProxy import com.rainy.token.data.repository.WorkBuddyRepository import dagger.Module import dagger.Provides @@ -223,6 +226,21 @@ object NetworkModule { balanceCache: BalanceCache ): Sub2ApiRepository = Sub2ApiRepository(okHttpClient, credentialRepository, balanceCache) + // ---- 本地反代网关 ---- + + @Provides + @Singleton + fun provideSub2ApiChatProxy( + okHttpClient: OkHttpClient, + credentialRepository: CredentialRepository + ): Sub2ApiChatProxy = Sub2ApiChatProxy(okHttpClient, credentialRepository) + + @Provides + @Singleton + fun provideLocalProxyServer( + sub2ApiChatProxy: Sub2ApiChatProxy + ): LocalProxyServer = KtorLocalProxyServer(sub2ApiChatProxy) + /** 余额缓存 DataStore(计划 7.1) */ @Provides @Singleton diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index 90a634c..cd10e79 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -15,6 +15,9 @@ navigationCompose = "2.9.8" # Network (Retrofit 2.11.0 配 OkHttp 4.12.0 — 稳定组合,避开 Retrofit 3.x + OkHttp 5.x 的前沿兼容性问题) retrofit = "2.11.0" + +# Ktor (本地反代 HTTP 服务) +ktor = "3.0.3" okhttp = "4.12.0" kotlinxSerializationJson = "1.7.3" retrofitKotlinxSerializationConverter = "1.0.0" @@ -67,6 +70,10 @@ androidx-navigation-compose = { group = "androidx.navigation", name = "navigatio retrofit = { group = "com.squareup.retrofit2", name = "retrofit", version.ref = "retrofit" } okhttp = { group = "com.squareup.okhttp3", name = "okhttp", version.ref = "okhttp" } okhttp-logging-interceptor = { group = "com.squareup.okhttp3", name = "logging-interceptor", version.ref = "okhttp" } +ktor-server-core = { group = "io.ktor", name = "ktor-server-core", version.ref = "ktor" } +ktor-server-cio = { group = "io.ktor", name = "ktor-server-cio", version.ref = "ktor" } +ktor-server-content-negotiation = { group = "io.ktor", name = "ktor-server-content-negotiation", version.ref = "ktor" } +ktor-serialization-kotlinx-json = { group = "io.ktor", name = "ktor-serialization-kotlinx-json", version.ref = "ktor" } kotlinx-serialization-json = { group = "org.jetbrains.kotlinx", name = "kotlinx-serialization-json", version.ref = "kotlinxSerializationJson" } retrofit-kotlinx-serialization-converter = { group = "com.jakewharton.retrofit", name = "retrofit2-kotlinx-serialization-converter", version.ref = "retrofitKotlinxSerializationConverter" }