* Create a pending write promise, optionally racing against a signal. * If the signal fires, the entry is removed from pendingWrites and the * promise rejects. Signal listeners are cleaned up on normal resolution.
(chunks: Uint8Array[], signal?: AbortSignal)
| 285 | * promise rejects. Signal listeners are cleaned up on normal resolution. |
| 286 | */ |
| 287 | private createPendingWrite(chunks: Uint8Array[], signal?: AbortSignal): Promise<void> { |
| 288 | return new Promise<void>((resolve, reject) => { |
| 289 | const entry: PendingWrite = { chunks, resolve, reject }; |
| 290 | this.pendingWrites.push(entry); |
| 291 | |
| 292 | if (!signal) return; |
| 293 | |
| 294 | const onAbort = () => { |
| 295 | // Remove from queue so it doesn't occupy a slot |
| 296 | const idx = this.pendingWrites.indexOf(entry); |
| 297 | if (idx !== -1) this.pendingWrites.removeAt(idx); |
| 298 | reject(signal.reason ?? new DOMException('Aborted', 'AbortError')); |
| 299 | }; |
| 300 | |
| 301 | // Wrap resolve/reject to clean up signal listener |
| 302 | const origResolve = entry.resolve; |
| 303 | const origReject = entry.reject; |
| 304 | entry.resolve = () => { |
| 305 | signal.removeEventListener('abort', onAbort); |
| 306 | origResolve(); |
| 307 | }; |
| 308 | entry.reject = (reason: Error) => { |
| 309 | signal.removeEventListener('abort', onAbort); |
| 310 | origReject(reason); |
| 311 | }; |
| 312 | |
| 313 | signal.addEventListener('abort', onAbort, { once: true }); |
| 314 | }); |
| 315 | } |
| 316 | |
| 317 | /** |
| 318 | * Signal end of stream. Returns total bytes written. |
no test coverage detected