Last active
September 23, 2026 19:34
-
-
Save mootari/d41dc0da35a82f501d8edb1a16ead7e6 to your computer and use it in GitHub Desktop.
Using linked lists for call throttling and pushable async queues.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| function deferred<T = void>() { | |
| let resolve: (v: T) => void; | |
| let reject: () => void; | |
| const promise = new Promise<T>((res, rej) => { | |
| resolve = res; | |
| reject = rej; | |
| }); | |
| return {promise, resolve: resolve!, reject: reject!}; | |
| } | |
| type Deferred<T> = ReturnType<typeof deferred<T>>; | |
| function asyncQueue<T>(): {queue: AsyncIterable<T, void, void>, push: (value: T) => void, end: () => void} { | |
| type Item = {pending: Deferred<IteratorResult<T, void>>, next?: Item}; | |
| let tail: Item = {pending: deferred()}; | |
| let head = tail; | |
| let drained = false; | |
| let closed = false; | |
| return { | |
| queue: { | |
| [Symbol.asyncIterator]: () => ({ | |
| next() { | |
| const p = head.pending.promise; | |
| if(head.next) head = head.next; | |
| else drained = true; | |
| return p; | |
| } | |
| }), | |
| }, | |
| async push(value: T) { | |
| if(closed) throw new Error("Cannot push to closed queue"); | |
| tail.pending.resolve({value, done: false}); | |
| tail = tail.next = {pending: deferred()}; | |
| if(drained) { | |
| head = head.next!; | |
| drained = false; | |
| } | |
| }, | |
| end() { | |
| closed = true; | |
| tail.pending.resolve({value: undefined, done: true}); | |
| } | |
| }; | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| function throttle(limit: number) { | |
| type Item = {run: () => void, next?: Item}; | |
| let tail: Item = {run: () => {}}; | |
| let head = tail; | |
| let added = 0, started = 0, done = 0; | |
| const next = () => { | |
| if(!head.next) return; | |
| head = head.next; | |
| head.run(); | |
| }; | |
| const enqueue = <T>(fn: () => T | Promise<T>) => new Promise<T>((resolve, reject) => { | |
| tail = tail.next = { | |
| async run() { | |
| ++started; | |
| await Promise.resolve().then(fn).then(resolve, reject); | |
| ++done; | |
| next(); | |
| }, | |
| }; | |
| ++added; | |
| if(started - done < limit) next(); | |
| }); | |
| return { | |
| enqueue, | |
| get waiting() { return added - started }, | |
| get running() { return started - done }, | |
| get done() { return done }, | |
| } as const; | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment