1/**
2 * implements a priority queue. jobs are be automatically cancelled via
3 * `AbortSignal` and rescheduled when higher priority tasks come in.
4 *
5 * @module
6 */
7
8/** schedule `item.run` to be called at some point on the global queue */
9export function run<T>(item: Job<T>): Promise<T> {
10 return global.run(item);
11}
12
13/** wrap a function into one that runs in the global queue */
14export function wrap<T, Args extends unknown[]>(
15 fn: (...args: Args) => Promise<T>,
16 getOpts?: (...args: Args) => Omit<Job<T>, "run">,
17): (...args: Args) => Promise<T> {
18 return global.wrap(fn, getOpts);
19}
20
21/** a definition for a queue job */
22export interface Job<out T> {
23 run: (signal: AbortSignal) => Promise<T>;
24 /** Approximately how many cores this job uses. */
25 cores?: number;
26 /** Higher priorities will cancel lower priority jobs. */
27 priority?: number;
28 /** If not `true`, the given AbortSignal will never abort. */
29 cancelable?: boolean;
30 /** @default false */
31 causeCancelations?: boolean;
32}
33
34export class PriorityQueue {
35 /** Read only */
36 concurrency: number;
37 /** Read only */
38 coresRemain: number;
39 /** Read only */
40 jobs: Array<Internal> = [];
41 /** Read only */
42 active: Array<Internal> = [];
43 /** Can be assigned to react to queue changes. */
44 onChange: (() => void) | null = null;
45
46 constructor(concurrency: number) {
47 this.concurrency = concurrency;
48 this.coresRemain = concurrency;
49 }
50
51 /** schedule `item.run` to be called at some point. */
52 run<T>(item: Job<T>): Promise<T> {
53 const { resolve, reject, promise } = Promise.withResolvers<T>();
54 const job = {
55 cores: 1,
56 cancelable: false,
57 priority: 1,
58 ...item,
59 ctrl: null,
60 cancelled: false,
61 resolve,
62 reject,
63 } satisfies Internal<T> as Internal;
64
65 if (this.canRun(job)) {
66 this.startJob(job);
67 } else {
68 if (job.causeCancelations ?? true) {
69 // Attempt to cancel other jobs to allow pushing this one up.
70 let remain = this.coresRemain;
71 for (
72 let i = this.active.length - 1;
73 i >= 0 && remain < job.cores;
74 i -= 1
75 ) {
76 const other = this.active[i]!;
77 if (other.priority > job.priority) break;
78 if (other.cancelable && !other.cancelled) {
79 other.cancelled = true;
80 UNWRAP(other.ctrl).abort(new RescheduleError());
81 remain += other.cores;
82 }
83 }
84 }
85
86 insertSorted(this.jobs, job);
87 }
88
89 return promise;
90 }
91
92 /** wrap a function into one that runs in the queue */
93 wrap<T, Args extends unknown[]>(
94 fn: (...args: Args) => Promise<T>,
95 getOpts?: (...args: Args) => Omit<Job<T>, "run">,
96 ): (...args: Args) => Promise<T> {
97 return (...args: Args) => this.run({ ...getOpts?.(...args), run: () => fn(...args) });
98 }
99
100 startJob(job: Internal) {
101 const { jobs, active } = this;
102 this.coresRemain -= job.cores;
103
104 // insert the job sorted
105 insertSorted(this.active, job);
106 this.onChange?.();
107
108 const { run, resolve, reject } = job;
109
110 const { signal } = (job.ctrl = new AbortController());
111 run(signal)
112 .finally(() => {
113 this.coresRemain += job.cores;
114 const index = active.indexOf(job);
115 if (index === -1) throw new Error("fyuc");
116 active.splice(index, 1);
117 if (this.coresRemain % 1 !== 0) {
118 this.coresRemain = active.reduce(
119 (acc, { cores }) => acc - cores,
120 this.concurrency,
121 );
122 }
123 })
124 .then(
125 (result) => {
126 resolve(result);
127 },
128 (error) => {
129 if (error instanceof RescheduleError) {
130 // handle re-scheduling
131 job.cancelled = false;
132 insertSorted(jobs, job);
133 } else {
134 reject(error);
135 }
136 },
137 )
138 .finally(() => {
139 // start next job
140 while (jobs[0]) {
141 if (this.canRun(jobs[0])) this.startJob(UNWRAP(jobs.shift()));
142 else break;
143 }
144 });
145 }
146
147 canRun({ cores }: Internal): boolean {
148 const { coresRemain, concurrency } = this;
149 return coresRemain >= cores
150 || (cores > concurrency && coresRemain === concurrency);
151 }
152}
153
154interface Internal<T = unknown> extends Job<unknown> {
155 cores: number;
156 priority: number;
157 cancelable: boolean;
158 ctrl: AbortController | null;
159 cancelled: boolean;
160 reject: (error: unknown) => void;
161 resolve: (value: T) => void;
162}
163
164export class RescheduleError extends Error {}
165
166function insertSorted<T extends { priority: number }>(arr: T[], item: T) {
167 for (let i = 0; i < arr.length; i += 1) {
168 if (item.priority > arr[i]!.priority) {
169 arr.splice(i, 0, item);
170 return;
171 }
172 }
173 arr.push(item);
174}
175
176declare var navigator: { hardwareConcurrency: number };
177const global = new PriorityQueue(navigator.hardwareConcurrency);
178
179/** adjust the global queue's concurrency */
180export function setConcurrency(n: number) {
181 ASSERT(Number.isInteger(n) && n > 0, `setConcurrency(${n})`);
182 global.coresRemain += n - global.concurrency;
183 global.concurrency = n;
184}
185
186import { ASSERT, UNWRAP } from "./assert.ts";