|
|
|
@ -20,11 +20,13 @@ import { useCallback, useEffect, useMemo, useRef, useState } from "react"; |
|
|
|
import { useStream } from "@langchain/react"; |
|
|
|
import { ProtocolSseTransportAdapter } from "@langchain/langgraph-sdk"; |
|
|
|
import { ROOT_PUMP_CHANNELS } from "@langchain/langgraph-sdk/stream"; |
|
|
|
import { getCompatibleToken } from "../../../utils/authCompat"; |
|
|
|
import { |
|
|
|
clearSuperAgentAccessToken, |
|
|
|
getSuperAgentFrameworkApiUrl, |
|
|
|
withAgentAuthorization, |
|
|
|
} from "../utils/agentApi"; |
|
|
|
import { getSuperAgentErrorMessage } from "../utils/agentApiParse"; |
|
|
|
import { superAgentLog, superAgentWarn } from "../utils/debugLog"; |
|
|
|
|
|
|
|
const BUSINESS_CHANNELS = [ |
|
|
|
"custom", |
|
|
|
@ -254,9 +256,16 @@ class SingleSubscriberSseTransportAdapter extends ProtocolSseTransportAdapter { |
|
|
|
async getState() { |
|
|
|
const requestVersion = ++this.stateRequestVersion; |
|
|
|
const requestedThreadId = String(this.threadId || "").trim(); |
|
|
|
// Agent 当前未提供 /threads/:threadId/state;useStream hydrate 只能使用
|
|
|
|
// 已从事件流获得的快照,不能为新会话触发一个必然失败的 HTTP 请求。
|
|
|
|
const state = this.latestState; |
|
|
|
if (!requestedThreadId) { |
|
|
|
this.latestState = null; |
|
|
|
this.onStateChange?.(); |
|
|
|
return null; |
|
|
|
} |
|
|
|
|
|
|
|
// ProtocolSseTransportAdapter 的 getState() 会请求当前 thread 的
|
|
|
|
// /state;其 fetchImpl 已绑定 fetchUseStreamTransport,因此会复用
|
|
|
|
// 8005 代理、Agent JWT 和 401 换票逻辑。
|
|
|
|
const state = await super.getState(); |
|
|
|
// 旧 state 快照不能在切换 thread 后覆盖新会话。
|
|
|
|
if ( |
|
|
|
requestVersion === this.stateRequestVersion && |
|
|
|
@ -534,14 +543,21 @@ const getEventDedupeKey = (event) => { |
|
|
|
|
|
|
|
const normalizeStreamError = (error) => { |
|
|
|
const status = error?.status || error?.response?.status; |
|
|
|
if (status === 401 || status === 403) return new Error("登录状态已失效,请重新登录"); |
|
|
|
if (status === 404) return new Error("智能体服务暂不可用"); |
|
|
|
const normalized = error instanceof Error ? error : new Error("智能体协议响应错误"); |
|
|
|
normalized.status = status || normalized.status || 0; |
|
|
|
normalized.code = error?.code || normalized.code || ""; |
|
|
|
normalized.requestId = error?.requestId || error?.request_id || normalized.requestId || ""; |
|
|
|
normalized.message = getSuperAgentErrorMessage(normalized, "智能体协议响应错误"); |
|
|
|
if (status === 404) normalized.message = "智能体服务暂不可用"; |
|
|
|
const message = String(error?.message || ""); |
|
|
|
if (status === 429 || /(?:^|\D)429(?:\D|$)|too many requests|rate limit/i.test(message)) { |
|
|
|
return new Error("事件流连接繁忙,正在自动恢复,请稍候"); |
|
|
|
normalized.message = "事件流连接繁忙,正在自动恢复,请稍候"; |
|
|
|
} |
|
|
|
if (/sse|stream|network|fetch/i.test(message)) return new Error("智能体流式连接异常"); |
|
|
|
return error instanceof Error ? error : new Error("智能体协议响应错误"); |
|
|
|
if (/sse|stream|network|fetch/i.test(message)) normalized.message = "智能体流式连接异常"; |
|
|
|
if (normalized.requestId && !normalized.message.includes(normalized.requestId)) { |
|
|
|
normalized.message = `${normalized.message}(request_id: ${normalized.requestId})`; |
|
|
|
} |
|
|
|
return normalized; |
|
|
|
}; |
|
|
|
|
|
|
|
const useSuperAgentStream = ({ |
|
|
|
@ -584,16 +600,31 @@ const useSuperAgentStream = ({ |
|
|
|
const localRunHasSeenLoadingRef = useRef(false); |
|
|
|
const activeSessionIdRef = useRef(""); |
|
|
|
activeSessionIdRef.current = activeSessionId; |
|
|
|
const loginToken = getCompatibleToken(); |
|
|
|
// 保持登录态变化时重建 transport;Agent JWT 在 fetch 拦截器中换票。
|
|
|
|
// 会话切换由 callerOptions 的会话级 fetch 变化触发新 controller。
|
|
|
|
// Agent 请求只使用短时 Agent JWT;ai-center token 只用于 query api 换票。
|
|
|
|
const defaultHeaders = useMemo( |
|
|
|
() => ({ |
|
|
|
Accept: "application/json, text/event-stream", |
|
|
|
...(loginToken ? { "X-PEP-Token": loginToken } : {}), |
|
|
|
}), |
|
|
|
[loginToken] |
|
|
|
[] |
|
|
|
); |
|
|
|
const fetchAuthorizedStream = useCallback(async (input, requestInit) => { |
|
|
|
let authorizedInit = await withAgentAuthorization(requestInit); |
|
|
|
let response = await fetchEventStreamWithRecovery(input, authorizedInit); |
|
|
|
if (response.status !== 401 || requestInit?.signal?.aborted) return response; |
|
|
|
|
|
|
|
superAgentWarn("stream.auth", "agent_token_refresh", { |
|
|
|
endpoint: String(input?.url || input || "").split("?")[0], |
|
|
|
status: response.status, |
|
|
|
}); |
|
|
|
clearSuperAgentAccessToken(); |
|
|
|
authorizedInit = await withAgentAuthorization(requestInit); |
|
|
|
response = await fetchEventStreamWithRecovery(input, authorizedInit); |
|
|
|
superAgentLog("stream.auth", "agent_token_refresh_result", { |
|
|
|
endpoint: String(input?.url || input || "").split("?")[0], |
|
|
|
status: response.status, |
|
|
|
}); |
|
|
|
return response; |
|
|
|
}, []); |
|
|
|
const fetchUseStreamTransport = useCallback(async (...args) => { |
|
|
|
const [input, init = {}] = args; |
|
|
|
if (isUseStreamHistoryDiscoveryRequest(input)) { |
|
|
|
@ -629,9 +660,8 @@ const useSuperAgentStream = ({ |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
const authorizedInit = await withAgentAuthorization(requestInit); |
|
|
|
return fetchEventStreamWithRecovery(input, authorizedInit); |
|
|
|
}, []); |
|
|
|
return fetchAuthorizedStream(input, requestInit); |
|
|
|
}, [fetchAuthorizedStream]); |
|
|
|
const eventTransport = useMemo( |
|
|
|
() => { |
|
|
|
const transport = new SingleSubscriberSseTransportAdapter({ |
|
|
|
@ -898,6 +928,11 @@ const useSuperAgentStream = ({ |
|
|
|
// 切换会话后旧 Thread 的监听器可能在清理前再收到一个事件,
|
|
|
|
// 不能把旧会话的事件写入新会话的流式消息状态。
|
|
|
|
if (activeSessionIdRef.current !== listenerSessionId) return; |
|
|
|
superAgentLog("stream.event", "received", { |
|
|
|
threadId: listenerSessionId, |
|
|
|
eventType: event?.method || "unknown", |
|
|
|
lifecycleType: event?.params?.data?.event || "", |
|
|
|
}); |
|
|
|
handleMessageEvent(event); |
|
|
|
const eventPayload = getEventPayload(event) || {}; |
|
|
|
if (shouldRefreshRuntimeState({ |
|
|
|
|