/* * AgentMail 鸿蒙客户端 — SSE 多账号实时推送服务 * 每个账号各建一条 SSE 连接(各带自己的 user_key) * AccountManager 维护连接集合,按 accountId 分发事件 */ import { http } from '@kit.NetworkKit'; import { BusinessError } from '@kit.BasicServicesKit'; import { hilog } from '@kit.PerformanceAnalysisKit'; /* UTF-8 解码用(见 `arrayBufferToString` 的注释:逐字节 fromCharCode 会把中文解成乱码) */ import { util } from '@kit.ArkTS'; import { AccountManager, AccountInfo } from './AccountManager'; import { normalizeApiBase } from '../model/ApiBase'; const DOMAIN = 0x0001; const TAG = 'SseService'; /** SSE 事件数据 */ export class SseEvent { type: string = ''; data: string = ''; accountId: string = ''; } /** SSE 回调 */ export type SseListener = (event: SseEvent) => void; /** SSE 连接状态 */ export type SseStatus = 'disconnected' | 'connecting' | 'connected'; /** SSE 状态回调 */ export type SseStatusListener = (status: SseStatus) => void; /** 单个账号的 SSE 连接状态(纯数据) */ export class SseConnection { accountId: string = ''; server: string = ''; token: string = ''; httpRequest: http.HttpRequest | null = null; listeners: SseListener[] = []; status: SseStatus = 'disconnected'; reconnectTimer: number = 0; buffer: string = ''; connected: boolean = false; /* * ★★ 2026-09-24:**每连接一个** 解码器(不要共享)。 * `decodeToString(..., {stream:true})` 会把「被 TCP 分片切断的半个汉字」 * 存在**该解码器实例内部**,等下一片到了再拼。多个连接共用同一个实例 * 就会把一个账号的半个字与另一个账号的半截拼在一起 ⇒ 乱码。 */ decoder: util.TextDecoder | null = null; } export class SseService { private static instance: SseService | null = null; private connections: Map = new Map(); private globalListeners: SseListener[] = []; private globalStatusListeners: SseStatusListener[] = []; static getInstance(): SseService { if (SseService.instance === null) { SseService.instance = new SseService(); } return SseService.instance; } private constructor() {} /** 添加全局事件监听(所有账号的事件都会收到) */ addListener(listener: SseListener): void { this.globalListeners.push(listener); } /** 移除全局事件监听 */ removeListener(listener: SseListener): void { const idx: number = this.globalListeners.indexOf(listener); if (idx >= 0) { this.globalListeners.splice(idx, 1); } } /** 添加全局状态监听 */ addStatusListener(listener: SseStatusListener): void { this.globalStatusListeners.push(listener); } /** 移除全局状态监听 */ removeStatusListener(listener: SseStatusListener): void { const idx: number = this.globalStatusListeners.indexOf(listener); if (idx >= 0) { this.globalStatusListeners.splice(idx, 1); } } /** 为指定账号建立 SSE 连接 */ connectForAccount(accountId: string, server: string, token: string): void { let conn: SseConnection | undefined = this.connections.get(accountId); if (conn !== undefined && conn.connected) { hilog.info(DOMAIN, TAG, 'account %{public}s already connected', accountId); return; } if (conn === undefined) { conn = new SseConnection(); conn.accountId = accountId; this.connections.set(accountId, conn); } conn.server = normalizeApiBase(server); conn.token = token; conn.connected = true; this.setConnStatus(conn, 'connecting'); this.doConnect(conn); } /** 断开指定账号的 SSE 连接 */ disconnectAccount(accountId: string): void { const conn: SseConnection | undefined = this.connections.get(accountId); if (conn === undefined) { return; } conn.connected = false; if (conn.reconnectTimer !== 0) { clearTimeout(conn.reconnectTimer); conn.reconnectTimer = 0; } if (conn.httpRequest !== null) { conn.httpRequest.destroy(); conn.httpRequest = null; } this.setConnStatus(conn, 'disconnected'); this.connections.delete(accountId); } /** 断开所有连接 */ disconnectAll(): void { const keys: string[] = []; this.connections.forEach((_conn: SseConnection, key: string) => { keys.push(key); }); for (let i = 0; i < keys.length; i++) { this.disconnectAccount(keys[i]); } } /** 根据 AccountManager 连接所有账号 */ async connectAll(acctMgr: AccountManager): Promise { await acctMgr.load(); const accounts: AccountInfo[] = acctMgr.getAccounts(); for (let i = 0; i < accounts.length; i++) { const acct: AccountInfo = accounts[i]; this.connectForAccount(acct.id, acct.server, acct.token); } } /** 获取指定账号的连接状态 */ getStatusForAccount(accountId: string): SseStatus { const conn: SseConnection | undefined = this.connections.get(accountId); if (conn === undefined) { return 'disconnected'; } return conn.status; } /** 给指定账号添加事件监听 */ addListenerForAccount(accountId: string, listener: SseListener): void { let conn: SseConnection | undefined = this.connections.get(accountId); if (conn === undefined) { conn = new SseConnection(); conn.accountId = accountId; this.connections.set(accountId, conn); } conn.listeners.push(listener); } /** 移除指定账号的事件监听 */ removeListenerForAccount(accountId: string, listener: SseListener): void { const conn: SseConnection | undefined = this.connections.get(accountId); if (conn !== undefined) { const idx: number = conn.listeners.indexOf(listener); if (idx >= 0) { conn.listeners.splice(idx, 1); } } } private setConnStatus(conn: SseConnection, status: SseStatus): void { if (conn.status !== status) { conn.status = status; hilog.info(DOMAIN, TAG, 'status[%{public}s]: %{public}s', conn.accountId, status); for (let i = 0; i < this.globalStatusListeners.length; i++) { this.globalStatusListeners[i](status); } } } private doConnect(conn: SseConnection): void { if (!conn.connected) { return; } const httpRequest = http.createHttp(); conn.httpRequest = httpRequest; /* 重连要换新的:旧解码器里可能还卡着上次断线时那半个字 */ conn.decoder = util.TextDecoder.create('utf-8', { ignoreBOM: true }); const url: string = conn.server + '/events/stream'; const header: Record = {}; if (conn.token.length > 0) { header['Authorization'] = 'Bearer ' + conn.token; } hilog.info(DOMAIN, TAG, 'connecting account %{public}s to %{public}s', conn.accountId, url); httpRequest.on('dataReceive', (data: ArrayBuffer) => { const text: string = this.arrayBufferToString(conn, data); conn.buffer += text; this.processBuffer(conn); }); httpRequest.on('dataEnd', () => { hilog.info(DOMAIN, TAG, 'dataEnd for account %{public}s', conn.accountId); this.setConnStatus(conn, 'disconnected'); this.scheduleReconnect(conn); }); httpRequest.on('headersReceive', (_headers: Object) => { hilog.info(DOMAIN, TAG, 'headers received for account %{public}s', conn.accountId); }); const options: http.HttpRequestOptions = { method: http.RequestMethod.GET, header: header, connectTimeout: 10000, readTimeout: 0 }; httpRequest.requestInStream(url, options, (err: BusinessError, code: number) => { if (err !== undefined && err !== null) { hilog.error(DOMAIN, TAG, 'requestInStream error account %{public}s: %{public}s', conn.accountId, err.message); this.setConnStatus(conn, 'disconnected'); this.scheduleReconnect(conn); return; } hilog.info(DOMAIN, TAG, 'requestInStream code=%{public}d account=%{public}s', code, conn.accountId); if (code === 200) { this.setConnStatus(conn, 'connected'); } else { hilog.error(DOMAIN, TAG, 'SSE failed code=%{public}d account=%{public}s', code, conn.accountId); this.setConnStatus(conn, 'disconnected'); this.scheduleReconnect(conn); } }); } private processBuffer(conn: SseConnection): void { const lines: string[] = conn.buffer.split('\n'); conn.buffer = lines.pop() ?? ''; let eventType: string = ''; let eventData: string = ''; for (let i = 0; i < lines.length; i++) { const line: string = lines[i]; if (line.length === 0) { if (eventType.length > 0 || eventData.length > 0) { const event: SseEvent = new SseEvent(); event.type = eventType.length > 0 ? eventType : 'message'; event.data = eventData; event.accountId = conn.accountId; this.dispatchEvent(conn, event); } eventType = ''; eventData = ''; } else if (line.startsWith('event:')) { eventType = line.substring(6).trim(); } else if (line.startsWith('data:')) { const newData: string = line.substring(5).trim(); if (eventData.length > 0) { eventData += '\n' + newData; } else { eventData = newData; } } } } private dispatchEvent(conn: SseConnection, event: SseEvent): void { hilog.info(DOMAIN, TAG, 'event[%{public}s]: %{public}s data: %{public}s', conn.accountId, event.type, event.data.substring(0, 100)); // 分发给账号级监听 for (let i = 0; i < conn.listeners.length; i++) { conn.listeners[i](event); } // 分发给全局监听 for (let i = 0; i < this.globalListeners.length; i++) { this.globalListeners[i](event); } } private scheduleReconnect(conn: SseConnection): void { if (!conn.connected) { return; } const timer: number | undefined = setTimeout(() => { conn.reconnectTimer = 0; this.doConnect(conn); }, 3000); conn.reconnectTimer = timer ?? 0; } private arrayBufferToString(conn: SseConnection, buffer: ArrayBuffer): string { /* * ★★ 2026-09-24 修:**UTF-8 解码**(用户:「鸿蒙 app 接收邮件的能力也有点不正常」)。 * * 原来这里是逐字节 `String.fromCharCode(uint8Array[i])` —— 那是 * **Latin-1 语义**:字节 ≥ 0x80 各自变成一个独立字符。 * SSE 流里一个汉字是 3 个 UTF-8 字节(如 「新」 = E6 96 B0) * ⇒ 解出来是 `æ\x96°` 这种乱码。 * * WebUI 没这个问题:它用浏览器原生 `EventSource`,规范本身按 UTF-8 * 解码(`client/electron/src/api/sse.ts:59`)。 * ⇒ 两端差在**解码这一步**,不是差在有没有 SSE。 * * ★ 为什么不用 `util.TextDecoder`:HarmonyOS 有 `@kit.ArkTS` 的 * `util.TextDecoder`(标准 API),比自己手写 UTF-8 状态机可靠得多 * —— 手写容易在「3 字节序列被 TCP 分片切成两段」时错(SSE 常见)。 * `stream: true` 正是为这种分段场景准备的:不完整的尾字节会 * 留在内部缓冲里,等下一片到了再拼。 */ const decoder: util.TextDecoder = conn.decoder !== null ? conn.decoder : util.TextDecoder.create('utf-8', { ignoreBOM: true }); conn.decoder = decoder; return decoder.decodeToString(new Uint8Array(buffer), { stream: true }); } }