MCPcopy Create free account
hub / github.com/WinterTC55/iter-streams / createPendingWrite

Method createPendingWrite

src/push.ts:287–315  ·  view source on GitHub ↗

* 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)

Source from the content-addressed store, hash-verified

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.

Callers 3

writeAsyncMethod · 0.95
broadcast.jsFile · 0.45
push.jsFile · 0.45

Calls 1

pushMethod · 0.65

Tested by

no test coverage detected