it-queueless-pushable
Advanced tools
| (function (root, factory) {(typeof module === 'object' && module.exports) ? module.exports = factory() : root.ItQueuelessPushable = factory()}(typeof self !== 'undefined' ? self : this, function () { | ||
| "use strict";var ItQueuelessPushable=(()=>{var o=Object.defineProperty;var d=Object.getOwnPropertyDescriptor;var c=Object.getOwnPropertyNames;var x=Object.prototype.hasOwnProperty;var f=(r,e)=>{for(var t in e)o(r,t,{get:e[t],enumerable:!0})},w=(r,e,t,n)=>{if(e&&typeof e=="object"||typeof e=="function")for(let s of c(e))!x.call(r,s)&&s!==t&&o(r,s,{get:()=>e[s],enumerable:!(n=d(e,s))||n.enumerable});return r};var v=r=>w(o({},"__esModule",{value:!0}),r);var N={};f(N,{queuelessPushable:()=>m});function i(){let r={};return r.promise=new Promise((e,t)=>{r.resolve=e,r.reject=t}),r}var a=class extends Error{type;code;constructor(e,t,n){super(e??"The operation was aborted"),this.type="aborted",this.name=n??"AbortError",this.code=t??"ABORT_ERR"}};async function h(r,e,t){if(e==null)return r;if(e.aborted)return r.catch(()=>{}),Promise.reject(new a(t?.errorMessage,t?.errorCode,t?.errorName));let n,s=new a(t?.errorMessage,t?.errorCode,t?.errorName);try{return await Promise.race([r,new Promise((p,l)=>{n=()=>{l(s)},e.addEventListener("abort",n)})])}finally{n!=null&&e.removeEventListener("abort",n)}}var u=class{readNext;haveNext;ended;nextResult;constructor(){this.ended=!1,this.readNext=i(),this.haveNext=i()}[Symbol.asyncIterator](){return this}async next(){if(this.nextResult==null&&await this.haveNext.promise,this.nextResult==null)throw new Error("HaveNext promise resolved but nextResult was undefined");let e=this.nextResult;return this.nextResult=void 0,this.readNext.resolve(),this.readNext=i(),e}async throw(e){return this.ended=!0,e!=null&&(this.haveNext.promise.catch(()=>{}),this.haveNext.reject(e)),{done:!0,value:void 0}}async return(){let e={done:!0,value:void 0};return this.ended=!0,this.nextResult=e,this.haveNext.resolve(),e}async push(e,t){await this._push(e,t)}async end(e,t){e!=null?await this.throw(e):await this._push(void 0,t)}async _push(e,t){if(e!=null&&this.ended)throw new Error("Cannot push value onto an ended pushable");for(;this.nextResult!=null;)await this.readNext.promise;e!=null?this.nextResult={done:!1,value:e}:(this.ended=!0,this.nextResult={done:!0,value:void 0}),this.haveNext.resolve(),this.haveNext=i(),await h(this.readNext.promise,t?.signal,t)}};function m(){return new u}return v(N);})(); | ||
| "use strict";var ItQueuelessPushable=(()=>{var o=Object.defineProperty;var d=Object.getOwnPropertyDescriptor;var c=Object.getOwnPropertyNames;var x=Object.prototype.hasOwnProperty;var f=(r,e)=>{for(var t in e)o(r,t,{get:e[t],enumerable:!0})},w=(r,e,t,n)=>{if(e&&typeof e=="object"||typeof e=="function")for(let s of c(e))!x.call(r,s)&&s!==t&&o(r,s,{get:()=>e[s],enumerable:!(n=d(e,s))||n.enumerable});return r};var v=r=>w(o({},"__esModule",{value:!0}),r);var N={};f(N,{queuelessPushable:()=>m});function i(){let r={};return r.promise=new Promise((e,t)=>{r.resolve=e,r.reject=t}),r}var a=class extends Error{type;code;constructor(e,t,n){super(e??"The operation was aborted"),this.type="aborted",this.name=n??"AbortError",this.code=t??"ABORT_ERR"}};async function h(r,e,t){if(e==null)return r;if(e.aborted)return r.catch(()=>{}),Promise.reject(new a(t?.errorMessage,t?.errorCode,t?.errorName));let n,s=new a(t?.errorMessage,t?.errorCode,t?.errorName);try{return await Promise.race([r,new Promise((p,l)=>{n=()=>{l(s)},e.addEventListener("abort",n)})])}finally{n!=null&&e.removeEventListener("abort",n)}}var u=class{readNext;haveNext;ended;nextResult;error;constructor(){this.ended=!1,this.readNext=i(),this.haveNext=i()}[Symbol.asyncIterator](){return this}async next(){if(this.nextResult==null&&await this.haveNext.promise,this.nextResult==null)throw new Error("HaveNext promise resolved but nextResult was undefined");let e=this.nextResult;return this.nextResult=void 0,this.readNext.resolve(),this.readNext=i(),e}async throw(e){return this.ended=!0,this.error=e,e!=null&&(this.haveNext.promise.catch(()=>{}),this.haveNext.reject(e)),{done:!0,value:void 0}}async return(){let e={done:!0,value:void 0};return this.ended=!0,this.nextResult=e,this.haveNext.resolve(),e}async push(e,t){await this._push(e,t)}async end(e,t){e!=null?await this.throw(e):await this._push(void 0,t)}async _push(e,t){if(e!=null&&this.ended)throw this.error??new Error("Cannot push value onto an ended pushable");for(;this.nextResult!=null;)await this.readNext.promise;e!=null?this.nextResult={done:!1,value:e}:(this.ended=!0,this.nextResult={done:!0,value:void 0}),this.haveNext.resolve(),this.haveNext=i(),await h(this.readNext.promise,t?.signal,t)}};function m(){return new u}return v(N);})(); | ||
| return ItQueuelessPushable})); |
@@ -31,5 +31,3 @@ /** | ||
| import { type RaceSignalOptions } from 'race-signal'; | ||
| export interface AbortOptions { | ||
| signal?: AbortSignal; | ||
| } | ||
| import type { AbortOptions } from 'abort-error'; | ||
| export interface Pushable<T> extends AsyncGenerator<T, void, unknown> { | ||
@@ -36,0 +34,0 @@ /** |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"index.d.ts","sourceRoot":"","sources":["../../src/index.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;;;;;;;;;GA4BG;AAGH,OAAO,EAAc,KAAK,iBAAiB,EAAE,MAAM,aAAa,CAAA;AAEhE,MAAM,WAAW,YAAY;IAC3B,MAAM,CAAC,EAAE,WAAW,CAAA;CACrB;AAED,MAAM,WAAW,QAAQ,CAAC,CAAC,CAAE,SAAQ,cAAc,CAAC,CAAC,EAAE,IAAI,EAAE,OAAO,CAAC;IACnE;;;;OAIG;IACH,GAAG,CAAC,GAAG,CAAC,EAAE,KAAK,EAAE,OAAO,CAAC,EAAE,YAAY,GAAG,iBAAiB,GAAG,OAAO,CAAC,IAAI,CAAC,CAAA;IAE3E;;;OAGG;IACH,IAAI,CAAC,KAAK,EAAE,CAAC,EAAE,OAAO,CAAC,EAAE,YAAY,GAAG,iBAAiB,GAAG,OAAO,CAAC,IAAI,CAAC,CAAA;CAC1E;AAoHD,wBAAgB,iBAAiB,CAAE,CAAC,KAAM,QAAQ,CAAC,CAAC,CAAC,CAEpD"} | ||
| {"version":3,"file":"index.d.ts","sourceRoot":"","sources":["../../src/index.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;;;;;;;;;GA4BG;AAGH,OAAO,EAAc,KAAK,iBAAiB,EAAE,MAAM,aAAa,CAAA;AAChE,OAAO,KAAK,EAAE,YAAY,EAAE,MAAM,aAAa,CAAA;AAE/C,MAAM,WAAW,QAAQ,CAAC,CAAC,CAAE,SAAQ,cAAc,CAAC,CAAC,EAAE,IAAI,EAAE,OAAO,CAAC;IACnE;;;;OAIG;IACH,GAAG,CAAC,GAAG,CAAC,EAAE,KAAK,EAAE,OAAO,CAAC,EAAE,YAAY,GAAG,iBAAiB,GAAG,OAAO,CAAC,IAAI,CAAC,CAAA;IAE3E;;;OAGG;IACH,IAAI,CAAC,KAAK,EAAE,CAAC,EAAE,OAAO,CAAC,EAAE,YAAY,GAAG,iBAAiB,GAAG,OAAO,CAAC,IAAI,CAAC,CAAA;CAC1E;AAsHD,wBAAgB,iBAAiB,CAAE,CAAC,KAAM,QAAQ,CAAC,CAAC,CAAC,CAEpD"} |
@@ -37,2 +37,3 @@ /** | ||
| nextResult; | ||
| error; | ||
| constructor() { | ||
@@ -63,2 +64,3 @@ this.ended = false; | ||
| this.ended = true; | ||
| this.error = err; | ||
| if (err != null) { | ||
@@ -101,3 +103,3 @@ // this can cause unhandled promise rejections if nothing is awaiting the | ||
| if (value != null && this.ended) { | ||
| throw new Error('Cannot push value onto an ended pushable'); | ||
| throw this.error ?? new Error('Cannot push value onto an ended pushable'); | ||
| } | ||
@@ -104,0 +106,0 @@ // wait for all values to be read |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"index.js","sourceRoot":"","sources":["../../src/index.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;;;;;;;;;GA4BG;AAEH,OAAO,QAAQ,EAAE,EAAwB,MAAM,SAAS,CAAA;AACxD,OAAO,EAAE,UAAU,EAA0B,MAAM,aAAa,CAAA;AAqBhE,MAAM,iBAAiB;IACb,QAAQ,CAAuB;IAC/B,QAAQ,CAAuB;IAC/B,KAAK,CAAS;IACd,UAAU,CAA+B;IAEjD;QACE,IAAI,CAAC,KAAK,GAAG,KAAK,CAAA;QAElB,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;QAC1B,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;IAC5B,CAAC;IAED,CAAC,MAAM,CAAC,aAAa,CAAC;QACpB,OAAO,IAAI,CAAA;IACb,CAAC;IAED,KAAK,CAAC,IAAI;QACR,IAAI,IAAI,CAAC,UAAU,IAAI,IAAI,EAAE,CAAC;YAC5B,wCAAwC;YACxC,MAAM,IAAI,CAAC,QAAQ,CAAC,OAAO,CAAA;QAC7B,CAAC;QAED,IAAI,IAAI,CAAC,UAAU,IAAI,IAAI,EAAE,CAAC;YAC5B,MAAM,IAAI,KAAK,CAAC,wDAAwD,CAAC,CAAA;QAC3E,CAAC;QAED,MAAM,UAAU,GAAG,IAAI,CAAC,UAAU,CAAA;QAClC,IAAI,CAAC,UAAU,GAAG,SAAS,CAAA;QAE3B,gDAAgD;QAChD,IAAI,CAAC,QAAQ,CAAC,OAAO,EAAE,CAAA;QACvB,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;QAE1B,OAAO,UAAU,CAAA;IACnB,CAAC;IAED,KAAK,CAAC,KAAK,CAAE,GAAW;QACtB,IAAI,CAAC,KAAK,GAAG,IAAI,CAAA;QAEjB,IAAI,GAAG,IAAI,IAAI,EAAE,CAAC;YAChB,yEAAyE;YACzE,6DAA6D;YAC7D,IAAI,CAAC,QAAQ,CAAC,OAAO,CAAC,KAAK,CAAC,GAAG,EAAE,GAAE,CAAC,CAAC,CAAA;YACrC,IAAI,CAAC,QAAQ,CAAC,MAAM,CAAC,GAAG,CAAC,CAAA;QAC3B,CAAC;QAED,MAAM,MAAM,GAAoC;YAC9C,IAAI,EAAE,IAAI;YACV,KAAK,EAAE,SAAS;SACjB,CAAA;QAED,OAAO,MAAM,CAAA;IACf,CAAC;IAED,KAAK,CAAC,MAAM;QACV,MAAM,MAAM,GAAoC;YAC9C,IAAI,EAAE,IAAI;YACV,KAAK,EAAE,SAAS;SACjB,CAAA;QAED,IAAI,CAAC,KAAK,GAAG,IAAI,CAAA;QACjB,IAAI,CAAC,UAAU,GAAG,MAAM,CAAA;QAExB,4CAA4C;QAC5C,IAAI,CAAC,QAAQ,CAAC,OAAO,EAAE,CAAA;QAEvB,OAAO,MAAM,CAAA;IACf,CAAC;IAED,KAAK,CAAC,IAAI,CAAE,KAAQ,EAAE,OAA0C;QAC9D,MAAM,IAAI,CAAC,KAAK,CAAC,KAAK,EAAE,OAAO,CAAC,CAAA;IAClC,CAAC;IAED,KAAK,CAAC,GAAG,CAAE,GAAW,EAAE,OAA0C;QAChE,IAAI,GAAG,IAAI,IAAI,EAAE,CAAC;YAChB,MAAM,IAAI,CAAC,KAAK,CAAC,GAAG,CAAC,CAAA;QACvB,CAAC;aAAM,CAAC;YACN,mBAAmB;YACnB,MAAM,IAAI,CAAC,KAAK,CAAC,SAAS,EAAE,OAAO,CAAC,CAAA;QACtC,CAAC;IACH,CAAC;IAEO,KAAK,CAAC,KAAK,CAAE,KAAS,EAAE,OAA0C;QACxE,IAAI,KAAK,IAAI,IAAI,IAAI,IAAI,CAAC,KAAK,EAAE,CAAC;YAChC,MAAM,IAAI,KAAK,CAAC,0CAA0C,CAAC,CAAA;QAC7D,CAAC;QAED,iCAAiC;QACjC,OAAO,IAAI,CAAC,UAAU,IAAI,IAAI,EAAE,CAAC;YAC/B,MAAM,IAAI,CAAC,QAAQ,CAAC,OAAO,CAAA;QAC7B,CAAC;QAED,IAAI,KAAK,IAAI,IAAI,EAAE,CAAC;YAClB,IAAI,CAAC,UAAU,GAAG,EAAE,IAAI,EAAE,KAAK,EAAE,KAAK,EAAE,CAAA;QAC1C,CAAC;aAAM,CAAC;YACN,IAAI,CAAC,KAAK,GAAG,IAAI,CAAA;YACjB,IAAI,CAAC,UAAU,GAAG,EAAE,IAAI,EAAE,IAAI,EAAE,KAAK,EAAE,SAAS,EAAE,CAAA;QACpD,CAAC;QAED,4CAA4C;QAC5C,IAAI,CAAC,QAAQ,CAAC,OAAO,EAAE,CAAA;QACvB,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;QAE1B,4EAA4E;QAC5E,6DAA6D;QAC7D,MAAM,UAAU,CACd,IAAI,CAAC,QAAQ,CAAC,OAAO,EACrB,OAAO,EAAE,MAAM,EACf,OAAO,CACR,CAAA;IACH,CAAC;CACF;AAED,MAAM,UAAU,iBAAiB;IAC/B,OAAO,IAAI,iBAAiB,EAAK,CAAA;AACnC,CAAC"} | ||
| {"version":3,"file":"index.js","sourceRoot":"","sources":["../../src/index.ts"],"names":[],"mappings":"AAAA;;;;;;;;;;;;;;;;;;;;;;;;;;;;GA4BG;AAEH,OAAO,QAAQ,EAAE,EAAwB,MAAM,SAAS,CAAA;AACxD,OAAO,EAAE,UAAU,EAA0B,MAAM,aAAa,CAAA;AAkBhE,MAAM,iBAAiB;IACb,QAAQ,CAAuB;IAC/B,QAAQ,CAAuB;IAC/B,KAAK,CAAS;IACd,UAAU,CAA+B;IACzC,KAAK,CAAQ;IAErB;QACE,IAAI,CAAC,KAAK,GAAG,KAAK,CAAA;QAElB,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;QAC1B,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;IAC5B,CAAC;IAED,CAAC,MAAM,CAAC,aAAa,CAAC;QACpB,OAAO,IAAI,CAAA;IACb,CAAC;IAED,KAAK,CAAC,IAAI;QACR,IAAI,IAAI,CAAC,UAAU,IAAI,IAAI,EAAE,CAAC;YAC5B,wCAAwC;YACxC,MAAM,IAAI,CAAC,QAAQ,CAAC,OAAO,CAAA;QAC7B,CAAC;QAED,IAAI,IAAI,CAAC,UAAU,IAAI,IAAI,EAAE,CAAC;YAC5B,MAAM,IAAI,KAAK,CAAC,wDAAwD,CAAC,CAAA;QAC3E,CAAC;QAED,MAAM,UAAU,GAAG,IAAI,CAAC,UAAU,CAAA;QAClC,IAAI,CAAC,UAAU,GAAG,SAAS,CAAA;QAE3B,gDAAgD;QAChD,IAAI,CAAC,QAAQ,CAAC,OAAO,EAAE,CAAA;QACvB,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;QAE1B,OAAO,UAAU,CAAA;IACnB,CAAC;IAED,KAAK,CAAC,KAAK,CAAE,GAAW;QACtB,IAAI,CAAC,KAAK,GAAG,IAAI,CAAA;QACjB,IAAI,CAAC,KAAK,GAAG,GAAG,CAAA;QAEhB,IAAI,GAAG,IAAI,IAAI,EAAE,CAAC;YAChB,yEAAyE;YACzE,6DAA6D;YAC7D,IAAI,CAAC,QAAQ,CAAC,OAAO,CAAC,KAAK,CAAC,GAAG,EAAE,GAAE,CAAC,CAAC,CAAA;YACrC,IAAI,CAAC,QAAQ,CAAC,MAAM,CAAC,GAAG,CAAC,CAAA;QAC3B,CAAC;QAED,MAAM,MAAM,GAAoC;YAC9C,IAAI,EAAE,IAAI;YACV,KAAK,EAAE,SAAS;SACjB,CAAA;QAED,OAAO,MAAM,CAAA;IACf,CAAC;IAED,KAAK,CAAC,MAAM;QACV,MAAM,MAAM,GAAoC;YAC9C,IAAI,EAAE,IAAI;YACV,KAAK,EAAE,SAAS;SACjB,CAAA;QAED,IAAI,CAAC,KAAK,GAAG,IAAI,CAAA;QACjB,IAAI,CAAC,UAAU,GAAG,MAAM,CAAA;QAExB,4CAA4C;QAC5C,IAAI,CAAC,QAAQ,CAAC,OAAO,EAAE,CAAA;QAEvB,OAAO,MAAM,CAAA;IACf,CAAC;IAED,KAAK,CAAC,IAAI,CAAE,KAAQ,EAAE,OAA0C;QAC9D,MAAM,IAAI,CAAC,KAAK,CAAC,KAAK,EAAE,OAAO,CAAC,CAAA;IAClC,CAAC;IAED,KAAK,CAAC,GAAG,CAAE,GAAW,EAAE,OAA0C;QAChE,IAAI,GAAG,IAAI,IAAI,EAAE,CAAC;YAChB,MAAM,IAAI,CAAC,KAAK,CAAC,GAAG,CAAC,CAAA;QACvB,CAAC;aAAM,CAAC;YACN,mBAAmB;YACnB,MAAM,IAAI,CAAC,KAAK,CAAC,SAAS,EAAE,OAAO,CAAC,CAAA;QACtC,CAAC;IACH,CAAC;IAEO,KAAK,CAAC,KAAK,CAAE,KAAS,EAAE,OAA0C;QACxE,IAAI,KAAK,IAAI,IAAI,IAAI,IAAI,CAAC,KAAK,EAAE,CAAC;YAChC,MAAM,IAAI,CAAC,KAAK,IAAI,IAAI,KAAK,CAAC,0CAA0C,CAAC,CAAA;QAC3E,CAAC;QAED,iCAAiC;QACjC,OAAO,IAAI,CAAC,UAAU,IAAI,IAAI,EAAE,CAAC;YAC/B,MAAM,IAAI,CAAC,QAAQ,CAAC,OAAO,CAAA;QAC7B,CAAC;QAED,IAAI,KAAK,IAAI,IAAI,EAAE,CAAC;YAClB,IAAI,CAAC,UAAU,GAAG,EAAE,IAAI,EAAE,KAAK,EAAE,KAAK,EAAE,CAAA;QAC1C,CAAC;aAAM,CAAC;YACN,IAAI,CAAC,KAAK,GAAG,IAAI,CAAA;YACjB,IAAI,CAAC,UAAU,GAAG,EAAE,IAAI,EAAE,IAAI,EAAE,KAAK,EAAE,SAAS,EAAE,CAAA;QACpD,CAAC;QAED,4CAA4C;QAC5C,IAAI,CAAC,QAAQ,CAAC,OAAO,EAAE,CAAA;QACvB,IAAI,CAAC,QAAQ,GAAG,QAAQ,EAAE,CAAA;QAE1B,4EAA4E;QAC5E,6DAA6D;QAC7D,MAAM,UAAU,CACd,IAAI,CAAC,QAAQ,CAAC,OAAO,EACrB,OAAO,EAAE,MAAM,EACf,OAAO,CACR,CAAA;IACH,CAAC;CACF;AAED,MAAM,UAAU,iBAAiB;IAC/B,OAAO,IAAI,iBAAiB,EAAK,CAAA;AACnC,CAAC"} |
+2
-1
| { | ||
| "name": "it-queueless-pushable", | ||
| "version": "1.0.2", | ||
| "version": "2.0.0", | ||
| "description": "A pushable queue that waits until a value is consumed before accepting another", | ||
@@ -139,2 +139,3 @@ "author": "Alex Potsides <alex@achingbrain.net>", | ||
| "dependencies": { | ||
| "abort-error": "^1.0.1", | ||
| "p-defer": "^4.0.1", | ||
@@ -141,0 +142,0 @@ "race-signal": "^1.1.3" |
+4
-5
@@ -33,7 +33,4 @@ /** | ||
| import { raceSignal, type RaceSignalOptions } from 'race-signal' | ||
| import type { AbortOptions } from 'abort-error' | ||
| export interface AbortOptions { | ||
| signal?: AbortSignal | ||
| } | ||
| export interface Pushable<T> extends AsyncGenerator<T, void, unknown> { | ||
@@ -59,2 +56,3 @@ /** | ||
| private nextResult: IteratorResult<T> | undefined | ||
| private error?: Error | ||
@@ -94,2 +92,3 @@ constructor () { | ||
| this.ended = true | ||
| this.error = err | ||
@@ -141,3 +140,3 @@ if (err != null) { | ||
| if (value != null && this.ended) { | ||
| throw new Error('Cannot push value onto an ended pushable') | ||
| throw this.error ?? new Error('Cannot push value onto an ended pushable') | ||
| } | ||
@@ -144,0 +143,0 @@ |
22409
0.95%3
50%+ Added
+ Added