refactor: 使用 fetch + ReadableStream 替代 EventSource 实现 SSE 连接,支持自定义 Authorization header
This commit is contained in:
@@ -41,8 +41,10 @@ function createSseConnection(options: UseSseOptions = {}) {
|
||||
const isConnected = computed(() => connectionState.value === SseConnectionState.CONNECTED);
|
||||
|
||||
let eventSource: EventSource | null = null;
|
||||
let abortController: AbortController | null = null;
|
||||
let isManualDisconnect = false;
|
||||
let connectionTimeoutTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
let reader: ReadableStreamDefaultReader<Uint8Array> | null = null;
|
||||
|
||||
const eventHandlers = new Map<string, Set<EventHandler>>();
|
||||
|
||||
@@ -101,7 +103,7 @@ function createSseConnection(options: UseSseOptions = {}) {
|
||||
const connect = () => {
|
||||
isManualDisconnect = false;
|
||||
|
||||
if (eventSource && connectionState.value === SseConnectionState.CONNECTED) {
|
||||
if (connectionState.value === SseConnectionState.CONNECTED) {
|
||||
log("SSE 已连接,跳过重复连接");
|
||||
return;
|
||||
}
|
||||
@@ -119,10 +121,8 @@ function createSseConnection(options: UseSseOptions = {}) {
|
||||
|
||||
connectionState.value = SseConnectionState.CONNECTING;
|
||||
|
||||
const separator = config.url.includes("?") ? "&" : "?";
|
||||
const fullUrl = `${config.url}${separator}token=${encodeURIComponent(token)}`;
|
||||
|
||||
eventSource = new EventSource(fullUrl);
|
||||
// 使用 fetch + ReadableStream 替代 EventSource,支持 Authorization header
|
||||
abortController = new AbortController();
|
||||
|
||||
connectionTimeoutTimer = setTimeout(() => {
|
||||
if (connectionState.value === SseConnectionState.CONNECTING) {
|
||||
@@ -131,16 +131,84 @@ function createSseConnection(options: UseSseOptions = {}) {
|
||||
}
|
||||
}, config.connectionTimeout);
|
||||
|
||||
eventSource.onopen = handleOpen;
|
||||
eventSource.onerror = handleError;
|
||||
eventSource.onmessage = handleMessage;
|
||||
fetch(config.url, {
|
||||
method: "GET",
|
||||
headers: {
|
||||
Authorization: `Bearer ${token}`,
|
||||
Accept: "text/event-stream",
|
||||
},
|
||||
signal: abortController.signal,
|
||||
})
|
||||
.then((response) => {
|
||||
if (!response.ok) {
|
||||
throw new Error(`HTTP ${response.status}`);
|
||||
}
|
||||
clearConnectionTimeout();
|
||||
connectionState.value = SseConnectionState.CONNECTED;
|
||||
log("SSE 连接已建立");
|
||||
return response.body?.getReader();
|
||||
})
|
||||
.then((r) => {
|
||||
if (!r) return;
|
||||
reader = r;
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = "";
|
||||
|
||||
// 注册所有已订阅的事件处理器
|
||||
eventHandlers.forEach((_, eventName) => {
|
||||
if (eventName !== "message") {
|
||||
eventSource!.addEventListener(eventName, handleCustomEvent(eventName) as EventListener);
|
||||
}
|
||||
});
|
||||
const processChunk = ({
|
||||
done,
|
||||
value,
|
||||
}: ReadableStreamReadResult<Uint8Array>): Promise<void> | void => {
|
||||
if (done) {
|
||||
connectionState.value = SseConnectionState.DISCONNECTED;
|
||||
log("SSE 连接已关闭");
|
||||
return;
|
||||
}
|
||||
|
||||
buffer += decoder.decode(value, { stream: true });
|
||||
const lines = buffer.split("\n");
|
||||
buffer = lines.pop() || ""; // 保留不完整的行
|
||||
|
||||
let currentEvent = "message";
|
||||
let currentData = "";
|
||||
|
||||
for (const line of lines) {
|
||||
if (line.startsWith(":")) continue; // 注释行(心跳)
|
||||
if (line.startsWith("event:")) {
|
||||
currentEvent = line.slice(6).trim();
|
||||
} else if (line.startsWith("data:")) {
|
||||
currentData = line.slice(5).trim();
|
||||
} else if (line === "") {
|
||||
// 空行表示事件结束
|
||||
if (currentData) {
|
||||
const handlers = eventHandlers.get(currentEvent);
|
||||
if (handlers) {
|
||||
try {
|
||||
const data = JSON.parse(currentData);
|
||||
handlers.forEach((h) => h(data));
|
||||
} catch {
|
||||
handlers.forEach((h) => h(currentData));
|
||||
}
|
||||
}
|
||||
log(`收到事件[${currentEvent}]:`, currentData);
|
||||
}
|
||||
currentEvent = "message";
|
||||
currentData = "";
|
||||
}
|
||||
}
|
||||
|
||||
return reader!.read().then(processChunk);
|
||||
};
|
||||
|
||||
return reader.read().then(processChunk);
|
||||
})
|
||||
.catch((err) => {
|
||||
if (err.name === "AbortError") {
|
||||
log("SSE 连接已主动断开");
|
||||
} else {
|
||||
logError("SSE 连接错误:", err);
|
||||
connectionState.value = SseConnectionState.DISCONNECTED;
|
||||
}
|
||||
});
|
||||
|
||||
log("正在建立 SSE 连接...");
|
||||
};
|
||||
@@ -151,10 +219,6 @@ function createSseConnection(options: UseSseOptions = {}) {
|
||||
}
|
||||
eventHandlers.get(eventName)!.add(handler);
|
||||
|
||||
if (eventName !== "message" && eventSource) {
|
||||
eventSource.addEventListener(eventName, handleCustomEvent(eventName) as EventListener);
|
||||
}
|
||||
|
||||
log(`已订阅事件: ${eventName}`);
|
||||
|
||||
return () => {
|
||||
@@ -173,13 +237,21 @@ function createSseConnection(options: UseSseOptions = {}) {
|
||||
isManualDisconnect = true;
|
||||
clearConnectionTimeout();
|
||||
|
||||
if (reader) {
|
||||
reader.cancel();
|
||||
reader = null;
|
||||
}
|
||||
if (abortController) {
|
||||
abortController.abort();
|
||||
abortController = null;
|
||||
}
|
||||
if (eventSource) {
|
||||
eventSource.close();
|
||||
eventSource = null;
|
||||
log("SSE 连接已断开");
|
||||
}
|
||||
|
||||
connectionState.value = SseConnectionState.DISCONNECTED;
|
||||
log("SSE 连接已断开");
|
||||
};
|
||||
|
||||
const cleanup = () => {
|
||||
|
||||
Reference in New Issue
Block a user