SmbStreamingHttpServer.ets 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503
  1. import Logger from '../util/Logger';
  2. import { fileIo } from '@kit.CoreFileKit';
  3. import { BusinessError } from '@kit.BasicServicesKit';
  4. import { httpServer, HttpRequest, HttpResponse } from '@webabcd/harmony-httpserver';
  5. import { readSmbRange, SmbConnectionInfo } from './SmbRangeReader';
  6. import { readFtpRange, FtpConnectionInfo } from './FtpRangeReader';
  7. import { connection } from '@kit.NetworkKit';
  8. const TAG = 'SmbStreamingHttpServer';
  9. const STREAM_ROUTE_PREFIX = '/smb-stream';
  10. const DEFAULT_HOST = '127.0.0.1';
  11. const MAX_CHUNK_SIZE = 512 * 1024; // 512KB per response
  12. const RANGE_WAIT_TIMEOUT_MS = 60000;
  13. const RANGE_POLL_INTERVAL_MS = 100;
  14. const SESSION_TTL_MS = 10 * 60 * 1000;
  15. const PORT_CANDIDATES: number[] = [18888, 18889, 18900, 18901, 18902];
  16. type StreamingProtocol = 'smb' | 'ftp';
  17. interface StreamingSessionOptions {
  18. sessionKey: string;
  19. cachePath: string;
  20. expectedSize?: number;
  21. mimeType?: string;
  22. fileName?: string;
  23. remotePath?: string;
  24. protocol?: StreamingProtocol;
  25. smbConnection?: SmbConnectionInfo;
  26. ftpConnection?: FtpConnectionInfo;
  27. }
  28. interface StreamingSession extends StreamingSessionOptions {
  29. lastAccess: number;
  30. }
  31. interface RangeInfo {
  32. start: number;
  33. end?: number;
  34. isRange: boolean;
  35. }
  36. interface ChunkResult {
  37. buffer: ArrayBuffer;
  38. start: number;
  39. end: number;
  40. total?: number;
  41. }
  42. class RangeNotSatisfiableError extends Error {
  43. availableRange: string;
  44. constructor(availableRange: string) {
  45. super('Requested Range Not Satisfiable');
  46. this.name = 'RangeNotSatisfiableError';
  47. this.availableRange = availableRange;
  48. }
  49. }
  50. export default class SmbStreamingHttpServer {
  51. private static instance?: SmbStreamingHttpServer;
  52. private port: number = -1;
  53. private sessions: Map<string, StreamingSession> = new Map();
  54. private startPromise?: Promise<void>;
  55. static getInstance(): SmbStreamingHttpServer {
  56. if (!SmbStreamingHttpServer.instance) {
  57. SmbStreamingHttpServer.instance = new SmbStreamingHttpServer();
  58. }
  59. return SmbStreamingHttpServer.instance;
  60. }
  61. async getStreamingUrl(options: StreamingSessionOptions): Promise<string> {
  62. await this.ensureServerStarted();
  63. const host = await this.resolveLocalHost();
  64. const now = Date.now();
  65. const existing = this.sessions.get(options.sessionKey);
  66. if (existing) {
  67. existing.lastAccess = now;
  68. existing.cachePath = options.cachePath;
  69. existing.expectedSize = options.expectedSize;
  70. existing.mimeType = options.mimeType;
  71. existing.fileName = options.fileName;
  72. existing.remotePath = options.remotePath;
  73. existing.protocol = options.protocol;
  74. existing.smbConnection = options.smbConnection;
  75. existing.ftpConnection = options.ftpConnection;
  76. Logger.debug(TAG, `更新SMB流会话: ${options.sessionKey}`);
  77. } else {
  78. this.sessions.set(options.sessionKey, {
  79. sessionKey: options.sessionKey,
  80. cachePath: options.cachePath,
  81. expectedSize: options.expectedSize,
  82. mimeType: options.mimeType,
  83. fileName: options.fileName,
  84. remotePath: options.remotePath,
  85. protocol: options.protocol,
  86. smbConnection: options.smbConnection,
  87. ftpConnection: options.ftpConnection,
  88. lastAccess: now
  89. });
  90. Logger.info(TAG, `注册SMB流会话: ${options.sessionKey}`);
  91. }
  92. this.cleanupExpiredSessions();
  93. return `http://${host}:${this.port}${STREAM_ROUTE_PREFIX}/${encodeURIComponent(options.sessionKey)}`;
  94. }
  95. private async resolveLocalHost(): Promise<string> {
  96. try {
  97. const netHandle = await connection.getDefaultNet();
  98. const addressGetter = (connection as unknown as { getAddressesByNetwork?: (handle: unknown) => Promise<unknown[]> })
  99. .getAddressesByNetwork;
  100. if (!addressGetter) {
  101. return DEFAULT_HOST;
  102. }
  103. const addresses = await addressGetter(netHandle);
  104. if (addresses && addresses.length > 0) {
  105. for (let i = 0; i < addresses.length; i++) {
  106. const raw = addresses[i] as Record<string, unknown>;
  107. const address = (raw.address ?? raw.addr ?? raw.ip) as string | undefined;
  108. if (!address) {
  109. continue;
  110. }
  111. if (address.startsWith('127.') || address === '0.0.0.0') {
  112. continue;
  113. }
  114. if (address.includes(':')) {
  115. continue;
  116. }
  117. return address;
  118. }
  119. }
  120. } catch (error) {
  121. Logger.warn(TAG, `解析本机IP失败: ${(error as Error).message}`);
  122. }
  123. return DEFAULT_HOST;
  124. }
  125. private async ensureServerStarted(): Promise<void> {
  126. if (this.port > 0) {
  127. return;
  128. }
  129. if (this.startPromise) {
  130. return this.startPromise;
  131. }
  132. this.startPromise = this.startServerInternal().finally(() => {
  133. this.startPromise = undefined;
  134. });
  135. return this.startPromise;
  136. }
  137. private async startServerInternal(): Promise<void> {
  138. httpServer.enableLog(false);
  139. this.port = await this.tryStartOnAvailablePort();
  140. httpServer.handleHttpRequestAsync((request: HttpRequest) => this.handleRequest(request));
  141. Logger.info(TAG, `SMB流HTTP服务启动,端口: ${this.port}`);
  142. }
  143. private async tryStartOnAvailablePort(): Promise<number> {
  144. for (let i = 0; i < PORT_CANDIDATES.length; i++) {
  145. const candidate = PORT_CANDIDATES[i];
  146. try {
  147. return await this.startOnPort(candidate);
  148. } catch (error) {
  149. const err = error as Error;
  150. Logger.warn(TAG, `端口${candidate}启动失败: ${err.message}`);
  151. }
  152. }
  153. throw new Error('无法启动SMB流HTTP服务,端口均不可用');
  154. }
  155. private startOnPort(port: number): Promise<number> {
  156. return new Promise((resolve, reject) => {
  157. try {
  158. httpServer.start(port, (error: BusinessError, realPort: number) => {
  159. if (error && error.code !== 0) {
  160. reject(new Error(`HTTP服务启动失败:${error.code}-${error.message}`));
  161. return;
  162. }
  163. resolve(realPort);
  164. });
  165. } catch (error) {
  166. reject(error as Error);
  167. }
  168. });
  169. }
  170. private async handleRequest(request: HttpRequest): Promise<HttpResponse> {
  171. try {
  172. this.cleanupExpiredSessions();
  173. if (!request || !request.url) {
  174. return this.buildPlainResponse(400, 'invalid request');
  175. }
  176. const sessionId = this.extractSessionId(request.url);
  177. Logger.debug(TAG, `收到请求 method=${request.method} url=${request.url} range=${this.getHeader(request.headers, 'range') ?? ''}`);
  178. if (!sessionId) {
  179. return this.buildPlainResponse(404, 'not found');
  180. }
  181. const session = this.sessions.get(sessionId);
  182. if (!session) {
  183. return this.buildPlainResponse(404, 'session expired');
  184. }
  185. session.lastAccess = Date.now();
  186. const method = (request.method || 'GET').toUpperCase();
  187. if (method !== 'GET' && method !== 'HEAD') {
  188. return this.buildPlainResponse(405, 'method not allowed');
  189. }
  190. if (method === 'HEAD') {
  191. return this.buildHeadResponse(session);
  192. }
  193. const rangeHeader = this.getHeader(request.headers, 'range');
  194. const rangeInfo = this.parseRangeHeader(rangeHeader);
  195. let chunk: ChunkResult;
  196. try {
  197. chunk = await this.readChunk(session, rangeInfo.start, rangeInfo.end);
  198. } catch (error) {
  199. const err = error as Error;
  200. if (err instanceof RangeNotSatisfiableError) {
  201. return this.buildRangeNotSatisfiableResponse(err.availableRange);
  202. }
  203. throw err;
  204. }
  205. return this.buildChunkResponse(session, chunk, rangeInfo.isRange);
  206. } catch (error) {
  207. const err = error as Error;
  208. Logger.error(TAG, `处理SMB流请求失败: ${err.message}`);
  209. return this.buildPlainResponse(500, err.message || 'internal error');
  210. }
  211. }
  212. private buildPlainResponse(statusCode: number, message: string): HttpResponse {
  213. return {
  214. statusCode,
  215. result: message,
  216. headers: {
  217. 'Content-Type': 'text/plain; charset=utf-8'
  218. }
  219. } as HttpResponse;
  220. }
  221. private buildRangeNotSatisfiableResponse(availableRange: string): HttpResponse {
  222. return {
  223. statusCode: 416,
  224. result: 'Requested Range Not Satisfiable',
  225. headers: {
  226. 'Content-Range': availableRange,
  227. 'Content-Type': 'text/plain; charset=utf-8'
  228. }
  229. } as HttpResponse;
  230. }
  231. private buildHeadResponse(session: StreamingSession): HttpResponse {
  232. const headers: Record<string, string> = {
  233. 'Accept-Ranges': 'bytes',
  234. 'Content-Type': session.mimeType || 'application/octet-stream'
  235. };
  236. if (session.expectedSize && session.expectedSize > 0) {
  237. headers['Content-Length'] = session.expectedSize.toString();
  238. }
  239. return {
  240. statusCode: 200,
  241. headers,
  242. result: ''
  243. } as HttpResponse;
  244. }
  245. private buildChunkResponse(session: StreamingSession, chunk: ChunkResult, isRange: boolean): HttpResponse {
  246. const headers: Record<string, string> = {
  247. 'Content-Type': session.mimeType || 'application/octet-stream',
  248. 'Accept-Ranges': 'bytes',
  249. 'Content-Length': (chunk.end - chunk.start + 1).toString()
  250. };
  251. if (chunk.total && chunk.total > 0) {
  252. headers['Content-Range'] = `bytes ${chunk.start}-${chunk.end}/${chunk.total}`;
  253. }
  254. return {
  255. statusCode: isRange ? 206 : 200,
  256. headers,
  257. result: chunk.buffer
  258. } as HttpResponse;
  259. }
  260. private extractSessionId(url: string): string | undefined {
  261. if (!url) {
  262. return undefined;
  263. }
  264. if (url.startsWith('http://') || url.startsWith('https://')) {
  265. const slashIndex = url.indexOf('/', url.indexOf('//') + 2);
  266. url = slashIndex >= 0 ? url.substring(slashIndex) : '/';
  267. }
  268. if (!url.startsWith(STREAM_ROUTE_PREFIX)) {
  269. return undefined;
  270. }
  271. let relative = url.substring(STREAM_ROUTE_PREFIX.length);
  272. if (relative.startsWith('/')) {
  273. relative = relative.substring(1);
  274. }
  275. const queryIndex = relative.indexOf('?');
  276. if (queryIndex >= 0) {
  277. relative = relative.substring(0, queryIndex);
  278. }
  279. if (!relative || relative.length === 0) {
  280. return undefined;
  281. }
  282. try {
  283. return decodeURIComponent(relative);
  284. } catch (error) {
  285. Logger.error(TAG, `sessionId解析失败: ${(error as Error).message}`);
  286. return undefined;
  287. }
  288. }
  289. private parseRangeHeader(header?: string): RangeInfo {
  290. if (!header || header.length === 0) {
  291. return { start: 0, isRange: false };
  292. }
  293. const match = header.match(/bytes=([0-9]*)-([0-9]*)/i);
  294. if (!match) {
  295. return { start: 0, isRange: false };
  296. }
  297. let start = match[1] ? parseInt(match[1]) : 0;
  298. const hasEnd = match[2] && match[2].length > 0;
  299. let end = hasEnd ? parseInt(match[2]) : undefined;
  300. if (isNaN(start) || start < 0) {
  301. start = 0;
  302. }
  303. if (end === undefined || isNaN(end) || end < start) {
  304. end = start + MAX_CHUNK_SIZE - 1;
  305. }
  306. const prefixRange = match[1] === undefined || match[1].length === 0;
  307. if (prefixRange) {
  308. const suffixBytes = parseInt(match[2]);
  309. if (!isNaN(suffixBytes) && suffixBytes > 0) {
  310. start = Math.max(0, end - suffixBytes + 1);
  311. }
  312. }
  313. if (end !== undefined && (isNaN(end) || end < start)) {
  314. end = start;
  315. }
  316. return { start, end, isRange: true };
  317. }
  318. private async readChunk(session: StreamingSession, requestedStart: number, requestedEnd?: number): Promise<ChunkResult> {
  319. const start = requestedStart >= 0 ? requestedStart : 0;
  320. let end = requestedEnd !== undefined ? requestedEnd : start + MAX_CHUNK_SIZE - 1;
  321. if (end - start + 1 > MAX_CHUNK_SIZE) {
  322. end = start + MAX_CHUNK_SIZE - 1;
  323. }
  324. const isFtp = session.protocol === 'ftp';
  325. const timeout = isFtp ? RANGE_WAIT_TIMEOUT_MS * 3 : RANGE_WAIT_TIMEOUT_MS;
  326. const pollInterval = isFtp ? RANGE_POLL_INTERVAL_MS * 2 : RANGE_POLL_INTERVAL_MS;
  327. const waitStart = Date.now();
  328. while (Date.now() - waitStart <= timeout) {
  329. const size = await this.tryGetFileSize(session.cachePath);
  330. Logger.debug(TAG, `range waiting start=${start} end=${end} currentSize=${size} cachePath=${session.cachePath} protocol=${session.protocol || 'unknown'}`);
  331. if (size >= 0 && size > start) {
  332. let fileEnd = Math.min(end, size - 1);
  333. if (session.expectedSize && session.expectedSize > 0) {
  334. fileEnd = Math.min(fileEnd, session.expectedSize - 1);
  335. }
  336. if (fileEnd >= start) {
  337. const length = fileEnd - start + 1;
  338. const buffer = new ArrayBuffer(length);
  339. const file = fileIo.openSync(session.cachePath, fileIo.OpenMode.READ_ONLY);
  340. try {
  341. fileIo.readSync(file.fd, buffer, { offset: start, length });
  342. } finally {
  343. fileIo.closeSync(file);
  344. }
  345. const total = session.expectedSize && session.expectedSize > 0 ? session.expectedSize : Math.max(size, fileEnd + 1);
  346. Logger.info(TAG, `从本地缓存返回范围 ${start}-${fileEnd},大小: ${length} bytes`);
  347. return {
  348. buffer,
  349. start,
  350. end: fileEnd,
  351. total
  352. };
  353. }
  354. }
  355. const hasLocalData = size > 0;
  356. const requestExceedsLocal = start >= size;
  357. if (requestExceedsLocal) {
  358. if (session.protocol === 'ftp' && !hasLocalData) {
  359. Logger.debug(TAG, 'FTP本地无数据,直接尝试远程读取');
  360. }
  361. const remoteChunk = await this.fetchRemoteChunk(session, start, end, size);
  362. if (remoteChunk && remoteChunk.buffer.byteLength > 0) {
  363. Logger.info(TAG, `远程读取成功,返回 ${remoteChunk.buffer.byteLength} bytes,实现边下边播`);
  364. return remoteChunk;
  365. }
  366. if (hasLocalData) {
  367. Logger.info(TAG, `用户请求未下载区域 ${start}-${end},已下载到 ${size - 1},返回HTTP 416让播放器跳回`);
  368. throw new RangeNotSatisfiableError(`bytes 0-${size - 1}/${session.expectedSize || '*'}`);
  369. }
  370. }
  371. if (session.expectedSize && session.expectedSize > 0 && start >= session.expectedSize) {
  372. throw new Error('请求范围超出文件大小');
  373. }
  374. await this.delay(pollInterval);
  375. }
  376. Logger.error(TAG, `range等待超时 start=${start} end=${end} cache=${session.cachePath} protocol=${session.protocol || 'unknown'}`);
  377. throw new Error(`等待${session.protocol === 'ftp' ? 'FTP' : 'SMB'}缓存数据超时`);
  378. }
  379. private async tryGetFileSize(path: string): Promise<number> {
  380. try {
  381. const stat = await fileIo.stat(path);
  382. return stat.size;
  383. } catch (error) {
  384. return -1;
  385. }
  386. }
  387. private getHeader(headers: Record<string, string> | undefined, name: string): string | undefined {
  388. if (!headers) {
  389. return undefined;
  390. }
  391. const target = name.toLowerCase();
  392. const keys = Object.keys(headers);
  393. for (let i = 0; i < keys.length; i++) {
  394. const key = keys[i];
  395. if (key.toLowerCase() === target) {
  396. return headers[key];
  397. }
  398. }
  399. return undefined;
  400. }
  401. private cleanupExpiredSessions(): void {
  402. const now = Date.now();
  403. const keys = Array.from(this.sessions.keys());
  404. for (let i = 0; i < keys.length; i++) {
  405. const key = keys[i];
  406. const session = this.sessions.get(key);
  407. if (session && now - session.lastAccess > SESSION_TTL_MS) {
  408. this.sessions.delete(key);
  409. Logger.info(TAG, `移除过期SMB流会话: ${key}`);
  410. }
  411. }
  412. }
  413. private delay(ms: number): Promise<void> {
  414. return new Promise(resolve => setTimeout(resolve, ms));
  415. }
  416. private async fetchRemoteChunk(
  417. session: StreamingSession,
  418. start: number,
  419. desiredEnd: number,
  420. currentLocalSize: number
  421. ): Promise<ChunkResult | undefined> {
  422. if (!session.protocol || !session.remotePath) {
  423. return undefined;
  424. }
  425. const length = desiredEnd - start + 1;
  426. if (length <= 0) {
  427. return undefined;
  428. }
  429. try {
  430. let data: ArrayBuffer | undefined;
  431. if (session.protocol === 'smb' && session.smbConnection) {
  432. data = readSmbRange(session.smbConnection, session.remotePath, start, length);
  433. } else if (session.protocol === 'ftp' && session.ftpConnection) {
  434. data = await readFtpRange(session.ftpConnection, session.remotePath, start, length);
  435. }
  436. if (!data || data.byteLength === 0) {
  437. return undefined;
  438. }
  439. const shouldWriteToCache = currentLocalSize >= 0 && start <= currentLocalSize;
  440. if (shouldWriteToCache) {
  441. await this.writeRangeToCache(session.cachePath, start, data);
  442. } else {
  443. Logger.debug(TAG, `跳过写入缓存,避免产生空洞: start=${start}, currentLocalSize=${currentLocalSize}`);
  444. }
  445. const chunkEnd = start + data.byteLength - 1;
  446. const totalSize = session.expectedSize && session.expectedSize > 0 ? session.expectedSize : Math.max(chunkEnd + 1, await this.tryGetFileSize(session.cachePath));
  447. return {
  448. buffer: data,
  449. start,
  450. end: chunkEnd,
  451. total: totalSize
  452. };
  453. } catch (error) {
  454. Logger.error(TAG, `远程读取失败: ${(error as Error).message}`);
  455. return undefined;
  456. }
  457. }
  458. private async writeRangeToCache(path: string, offset: number, data: ArrayBuffer): Promise<void> {
  459. const file = fileIo.openSync(path, fileIo.OpenMode.CREATE | fileIo.OpenMode.READ_WRITE);
  460. try {
  461. fileIo.writeSync(file.fd, data, { offset, length: data.byteLength });
  462. } finally {
  463. fileIo.closeSync(file);
  464. }
  465. }
  466. }