feat(proxy): M1a 本地 HTTP 服务 + Sub2API 透传

- 引入 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 注入
This commit is contained in:
Liuxinyu176 2026-10-08 22:58:33 +08:00
parent 9d3c004617
commit 85c6cf4283
5 changed files with 243 additions and 0 deletions

View File

@ -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)

View File

@ -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<Boolean> = _isRunning.asStateFlow()
private var server: EmbeddedServer<*, *>? = null
override fun start(config: ProxyServerConfig): Result<Unit> = 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)
}

View File

@ -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,
)

View File

@ -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

View File

@ -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" }