|
|
|
@ -22,7 +22,8 @@ import java.util.concurrent.atomic.AtomicBoolean |
|
|
|
*/ |
|
|
|
class CustomSseClientTransport( |
|
|
|
private val client: HttpClient, |
|
|
|
private val urlString: String?, |
|
|
|
private val serviceidString: String?, |
|
|
|
public val urlString: String?, |
|
|
|
private val reconnectionTime: Duration? = null, |
|
|
|
private val requestBuilder: HttpRequestBuilder.() -> Unit = {}, |
|
|
|
private val onConnectionLost: (() -> Unit)? = null |
|
|
|
@ -88,7 +89,7 @@ class CustomSseClientTransport( |
|
|
|
|
|
|
|
Triple(hostUrl, path, params) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "解析URL失败: $url, ${e.message}") |
|
|
|
Log.e(TAG, "$serviceidString 解析URL失败: $url, ${e.message}") |
|
|
|
Triple(url, "", emptyMap()) |
|
|
|
} |
|
|
|
} |
|
|
|
@ -102,9 +103,9 @@ class CustomSseClientTransport( |
|
|
|
session.incoming.collect { event -> |
|
|
|
when (event.event) { |
|
|
|
"error" -> { |
|
|
|
Log.e(TAG, "SSE错误: ${event.data}") |
|
|
|
Log.e(TAG, "$serviceidString SSE错误: ${event.data}") |
|
|
|
isConnected.set(false) |
|
|
|
val exception = Exception("SSE Error: ${event.data}") |
|
|
|
val exception = Exception("$serviceidString SSE Error: ${event.data}") |
|
|
|
_onError(exception) |
|
|
|
onConnectionLost?.invoke() |
|
|
|
throw exception |
|
|
|
@ -112,7 +113,7 @@ class CustomSseClientTransport( |
|
|
|
|
|
|
|
"open" -> { |
|
|
|
// SSE连接已打开 |
|
|
|
Log.d(TAG, "SSE连接已打开") |
|
|
|
Log.d(TAG, "$serviceidString SSE连接已打开") |
|
|
|
isConnected.set(true) |
|
|
|
} |
|
|
|
"ping" -> { |
|
|
|
@ -146,7 +147,7 @@ class CustomSseClientTransport( |
|
|
|
|
|
|
|
endpoint.complete(endpointWithParams) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "处理endpoint事件失败: ${e.message}", e) |
|
|
|
Log.e(TAG, "$serviceidString 处理endpoint事件失败: ${e.message}", e) |
|
|
|
_onError(e) |
|
|
|
close() |
|
|
|
error(e) |
|
|
|
@ -161,22 +162,22 @@ class CustomSseClientTransport( |
|
|
|
val message = json.decodeFromString<JSONRPCMessage>(data) |
|
|
|
_onMessage(message) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "解析JSON-RPC消息失败: ${e.message}", e) |
|
|
|
Log.e(TAG, "$serviceidString 解析JSON-RPC消息失败: ${e.message}", e) |
|
|
|
_onError(e) |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "处理事件失败: ${e.message}", e) |
|
|
|
Log.e(TAG, "$serviceidString 处理事件失败: ${e.message}", e) |
|
|
|
_onError(e) |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (e: CancellationException) { |
|
|
|
Log.d(TAG, "SSE事件收集被取消") |
|
|
|
Log.d(TAG, "$serviceidString SSE事件收集被取消") |
|
|
|
throw e |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "SSE连接异常断开: ${e.message}", e) |
|
|
|
Log.e(TAG, "$serviceidString SSE连接异常断开: ${e.message}", e) |
|
|
|
isConnected.set(false) |
|
|
|
_onError(e) |
|
|
|
onConnectionLost?.invoke() |
|
|
|
@ -199,13 +200,13 @@ class CustomSseClientTransport( |
|
|
|
|
|
|
|
// 检查session是否仍然活跃 |
|
|
|
if (session.coroutineContext[Job]?.isCancelled == true) { |
|
|
|
Log.w(TAG, "检测到SSE会话已取消") |
|
|
|
Log.w(TAG, "$serviceidString 检测到SSE会话已取消") |
|
|
|
isConnected.set(false) |
|
|
|
onConnectionLost?.invoke() |
|
|
|
break |
|
|
|
} |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "连接监控异常: ${e.message}", e) |
|
|
|
Log.e(TAG, "$serviceidString 连接监控异常: ${e.message}", e) |
|
|
|
isConnected.set(false) |
|
|
|
onConnectionLost?.invoke() |
|
|
|
break |
|
|
|
@ -219,7 +220,7 @@ class CustomSseClientTransport( |
|
|
|
*/ |
|
|
|
override suspend fun start() { |
|
|
|
if (!initialized.compareAndSet(false, true)) { |
|
|
|
Log.e(TAG, "传输层已经启动,不能重复启动") |
|
|
|
Log.e(TAG, "$serviceidString 传输层已经启动,不能重复启动") |
|
|
|
error("CustomSseClientTransport already started!") |
|
|
|
} |
|
|
|
|
|
|
|
@ -283,12 +284,30 @@ class CustomSseClientTransport( |
|
|
|
setBody(jsonString) |
|
|
|
} |
|
|
|
|
|
|
|
if (!response.status.isSuccess()) { |
|
|
|
val text = response.bodyAsText() |
|
|
|
// Log.e(TAG, "发送消息失败:URL:${urlString} HTTP ${response.status}, $text") |
|
|
|
// 优化处理百度API的202状态码 |
|
|
|
when { |
|
|
|
response.status.isSuccess() -> { |
|
|
|
// 2xx状态码都视为成功 |
|
|
|
Log.d(TAG, "$serviceidString 消息发送成功: HTTP ${response.status}") |
|
|
|
} |
|
|
|
response.status == HttpStatusCode.Accepted -> { |
|
|
|
// 202 Accepted - 百度API异步处理中,这是正常状态 |
|
|
|
Log.d(TAG, "$serviceidString 消息已被接受,正在异步处理: HTTP ${response.status}") |
|
|
|
} |
|
|
|
else -> { |
|
|
|
// 其他非成功状态码才记录为错误 |
|
|
|
val text = response.bodyAsText() |
|
|
|
Log.w(TAG, "$serviceidString 发送消息收到非成功状态码: URL:${messageEndpoint} HTTP ${response.status}, $text") |
|
|
|
|
|
|
|
// 根据具体状态码决定是否抛出异常 |
|
|
|
if (response.status.value >= 400) { |
|
|
|
// 4xx和5xx错误才抛出异常 |
|
|
|
throw Exception("HTTP ${response.status}: $text") |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "发送消息异常: ${e.message}", e) |
|
|
|
Log.e(TAG, "$serviceidString 发送消息异常: ${e.message}", e) |
|
|
|
_onError(e) |
|
|
|
throw e |
|
|
|
} |
|
|
|
@ -306,7 +325,7 @@ class CustomSseClientTransport( |
|
|
|
*/ |
|
|
|
override suspend fun close() { |
|
|
|
if (!initialized.get()) { |
|
|
|
Log.e(TAG, "关闭失败: 传输层未初始化") |
|
|
|
Log.e(TAG, "$serviceidString 关闭失败: 传输层未初始化") |
|
|
|
error("CustomSseClientTransport is not initialized!") |
|
|
|
} |
|
|
|
|
|
|
|
@ -317,6 +336,6 @@ class CustomSseClientTransport( |
|
|
|
job?.cancelAndJoin() |
|
|
|
connectionMonitorJob?.cancelAndJoin() |
|
|
|
|
|
|
|
Log.d(TAG, "CustomSseClientTransport已关闭") |
|
|
|
Log.d(TAG, "$serviceidString CustomSseClientTransport已关闭") |
|
|
|
} |
|
|
|
} |