authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2025-09-27 04:29:21-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2025-10-14 02:40:47-07:00
log6e8e86f58f873d0aa6151ac60aef9a637b531f07
treea6f17dfb2cce21e55bade71e1b19447aca9d7918
parent33bef63518232569f6494e9c8a4c48c9db3ec683
signature Commit is signed but in an unrecognized format.

feat(lib): add priority queue ported from work


1 files changed, 168 insertions(+), 0 deletions(-)

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