From 6e8e86f58f873d0aa6151ac60aef9a637b531f07 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Sat, 27 Sep 2025 04:29:21 -0700 Subject: [PATCH] feat(lib): add priority queue ported from work --- lib/queue.ts | 168 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 168 insertions(+) create mode 100644 lib/queue.ts diff --git a/lib/queue.ts b/lib/queue.ts new file mode 100644 index 0000000000000000000000000000000000000000..d90754bc8cf082b7dbfe33296092291503338bd7 --- /dev/null +++ b/lib/queue.ts @@ -0,0 +1,168 @@ +// Implements a priority queue. Jobs can be automatically cancelled via +// AbortSignal rescheduled if higher priority tasks come in. + +/** schedule `item.run` to be called at some point on the global queue */ +export function run(item: Job): Promise { + return global.run(item); +} + +/** wrap a function into one that runs in the global queue */ +export function wrap( + fn: (...args: Args) => Promise, + getOpts?: (...args: Args) => Omit, "run">, +): (...args: Args) => Promise { + return global.wrap(fn, getOpts); +} + +/** a definition for a queue job */ +export interface Job { + run: (signal: AbortSignal) => Promise; + /** Approximately how many cores this job uses. */ + cores?: number; + /** Higher priorities will cancel lower priority jobs. */ + priority?: number; + /** If not `true`, the given AbortSignal will never abort. */ + cancelable?: boolean; + /** @default false */ + causeCancelations?: boolean; +} + +export class PriorityQueue { + concurrency: number; + coresRemain: number; + jobs: Array = []; + active: Array = []; + onChange: (() => void) | null = null; + + constructor(concurrency: number) { + this.concurrency = concurrency; + this.coresRemain = concurrency; + } + + /** schedule `item.run` to be called at some point. */ + run(item: Job): Promise { + const { resolve, reject, promise } = Promise.withResolvers(); + const job = { + cores: 1, + cancelable: false, + priority: 1, + ...item, + ctrl: null, + cancelled: false, + resolve, + reject, + } satisfies Internal as Internal; + + if (this.canRun(job)) { + this.startJob(job); + } else { + if (job.causeCancelations ?? true) { + // Attempt to cancel other jobs to allow pushing this one up. + let remain = this.coresRemain; + for ( + let i = this.active.length - 1; + i >= 0 && remain < job.cores; + i -= 1 + ) { + const other = this.active[i]!; + if (other.priority > job.priority) break; + if (other.cancelable && !other.cancelled) { + other.cancelled = true; + UNWRAP(other.ctrl).abort(new RescheduleError()); + remain += other.cores; + } + } + } + + insertSorted(this.jobs, job); + } + + return promise; + } + + /** wrap a function into one that runs in the queue */ + wrap( + fn: (...args: Args) => Promise, + getOpts?: (...args: Args) => Omit, "run">, + ): (...args: Args) => Promise { + return (...args: Args) => + this.run({ ...getOpts?.(...args), run: () => fn(...args) }); + } + + startJob(job: Internal) { + const { jobs, active } = this; + this.coresRemain -= job.cores; + + // insert the job sorted + insertSorted(this.active, job); + this.onChange?.(); + + const { run, resolve, reject } = job; + + const { signal } = (job.ctrl = new AbortController()); + run(signal) + .finally(() => { + this.coresRemain += job.cores; + const index = active.indexOf(job); + if (index === -1) throw new Error("fyuc"); + active.splice(index, 1); + if (this.coresRemain % 1 !== 0) { + this.coresRemain = active.reduce( + (acc, { cores }) => acc - cores, + this.concurrency, + ); + } + }) + .then( + (result) => { + resolve(result); + }, + (error) => { + if (error instanceof RescheduleError) { + // handle re-scheduling + job.cancelled = false; + insertSorted(jobs, job); + } else { + reject(error); + } + }, + ) + .finally(() => { + // start next job + while (jobs[0]) { + if (this.canRun(jobs[0])) this.startJob(UNWRAP(jobs.shift())); + else break; + } + }); + } + + canRun({ cores }: Internal) { + const { coresRemain, concurrency } = this; + return coresRemain >= cores || + (cores > concurrency && coresRemain === concurrency); + } +} + +interface Internal extends Job { + cores: number; + priority: number; + cancelable: boolean; + ctrl: AbortController | null; + cancelled: boolean; + reject: (error: unknown) => void; + resolve: (value: T) => void; +} + +export class RescheduleError extends Error {} + +function insertSorted(arr: T[], item: T) { + for (let i = 0; i < arr.length; i += 1) { + if (item.priority > arr[i]!.priority) { + arr.splice(i, 0, item); + return; + } + } + arr.push(item); +} + +const global = new PriorityQueue(navigator.hardwareConcurrency); -- 2.54.0