1/**
2 * helpers to deal with promises and asynchronous execution.
3 *
4 * @module
5 */
6
7/*** @deprecated */
8interface ARCEValue<T> {
9 value: T;
10 [Symbol.dispose]: () => void;
11}
12
13/*** @deprecated */
14export function RefCountedExpirable<T>(
15 init: () => Promise<T>,
16 deinit: (value: T) => void,
17 expire: number = 5 * 60 * 1000,
18): () => Promise<ARCEValue<T>> {
19 let refs = 0;
20 let item: ARCEValue<T> | null = null;
21 let loading: Promise<ARCEValue<T>> | null = null;
22 let timer: ReturnType<typeof setTimeout> | null = null;
23
24 function deref() {
25 ASSERT(item !== null);
26 if (--refs !== 0) return;
27 ASSERT(timer === null);
28 timer = setTimeout(() => {
29 ASSERT(refs === 0);
30 ASSERT(loading === null);
31 ASSERT(item !== null);
32 deinit(item.value);
33 item = null;
34 timer = null;
35 }, expire);
36 }
37
38 return async function(): Promise<ARCEValue<T>> {
39 if (timer !== null) {
40 clearTimeout(timer);
41 timer = null;
42 }
43 if (item !== null) {
44 refs++;
45 return item;
46 }
47 if (loading !== null) {
48 refs++;
49 return loading;
50 }
51 const p = Promise.withResolvers<ARCEValue<T>>();
52 loading = p.promise;
53 try {
54 const value = await init();
55 item = { value, [Symbol.dispose]: deref };
56 refs++;
57 p.resolve(item);
58 return item;
59 } catch (e) {
60 p.reject(e);
61 throw e;
62 } finally {
63 loading = null;
64 }
65 };
66}
67
68/** evaluates `fn` one time, caching the result forever. */
69export function once<T>(fn: () => Promise<T>): () => Promise<T> {
70 let result: T | Promise<T> | null = null;
71 return async () => {
72 if (result) return result;
73 try {
74 result = await (result = fn());
75 } catch (err) {
76 result = null;
77 throw err;
78 }
79 return result;
80 };
81}
82
83/** a variable that can be watched for changes */
84export class Watch<T> {
85 #value: T;
86 #observers = new Set<WatchCallback<T>>();
87 constructor(value: T) {
88 this.#value = value;
89 }
90
91 on(cb: WatchCallback<T>): () => void {
92 this.#observers.add(cb);
93 return () => void this.#observers.delete(cb);
94 }
95 once(cb: WatchCallback<T>): void {
96 const once = (n: T, p: T) => (cb(n, p), this.#observers.delete(once));
97 this.#observers.add(once);
98 }
99
100 next(): Promise<T> {
101 return new Promise((resolve) => this.once(resolve));
102 }
103
104 until(condition: (v: T) => boolean): Promise<T> {
105 if (condition(this.#value)) return Promise.resolve(this.#value);
106 return new Promise<T>((resolve) => {
107 function check(next: T) {
108 if (condition(next)) resolve(next), release();
109 }
110 const release = this.on(check);
111 });
112 }
113
114 get value(): T {
115 return this.#value;
116 }
117 set value(next: T) {
118 const prev = this.#value;
119 if (prev !== next) {
120 this.#value = next;
121 this.#observers.forEach((cb) => cb(next, prev));
122 }
123 }
124}
125export type WatchCallback<T> = (next: T, prev: T) => void;
126
127/**
128 * when two requests with the same args come in at the same time, the response
129 * from them is shared. after returning, the data is no longer stored.
130 */
131export class DedupeConcurrent<Args extends unknown[], Ret> {
132 pending: Map<unknown, Promise<Ret>> = new Map();
133 fn: (...args: Args) => Promise<Ret>;
134 keyFn: (args: Args) => unknown = JSON.stringify;
135
136 constructor(
137 fn: (...args: Args) => Promise<Ret>,
138 keyFn: (args: Args) => unknown = JSON.stringify,
139 ) {
140 this.fn = fn;
141 this.keyFn = keyFn;
142 }
143
144 isPending(...args: Args): boolean {
145 return this.pending.has(this.keyFn(args));
146 }
147
148 /** Unbound to allow easily exporting this callback */
149 get = (...args: Args): Promise<Ret> => {
150 const k = this.keyFn(args);
151 let promise = this.pending.get(k);
152 if (promise) return promise;
153 promise = this.fn(...args);
154 this.pending.set(k, promise);
155 void promise.finally(() => void this.pending.delete(k));
156 return promise;
157 };
158}
159
160/**
161 * when two requests with the same args come, the result
162 * of the first is memoized.
163 */
164export class OnceMap<Args extends unknown[], Ret> {
165 cache: Map<unknown, Promise<Ret>> = new Map();
166 fn: (...args: Args) => Promise<Ret>;
167 keyFn: (args: Args) => unknown = JSON.stringify;
168
169 constructor(
170 fn: (...args: Args) => Promise<Ret>,
171 keyFn: (args: Args) => unknown = JSON.stringify,
172 ) {
173 this.fn = fn;
174 this.keyFn = keyFn;
175 }
176
177 has(...args: Args): boolean {
178 return this.cache.has(this.keyFn(args));
179 }
180
181 /** Unbound to allow easily exporting this callback */
182 getOrRun = (...args: Args): Promise<Ret> => {
183 const k = this.keyFn(args);
184 let promise = this.cache.get(k);
185 if (promise) return promise;
186 promise = this.fn(...args);
187 this.cache.set(k, promise);
188 return promise;
189 };
190}
191
192export class PromiseAggregator {
193 failures: unknown[] = [];
194 promises: Set<Promise<unknown>> = new Set();
195
196 push(p: Promise<unknown>) {
197 p = p.then(
198 () => {
199 this.promises.delete(p);
200 },
201 (err) => {
202 this.failures.push(err);
203 },
204 );
205 this.promises.add(p);
206 }
207
208 async all() {
209 while (this.promises.size > 0) {
210 const wait = [...this.promises];
211 this.promises.clear();
212 await Promise.all(wait);
213 }
214 if (this.failures.length > 0) {
215 const agg = new AggregateError(this.failures);
216 this.failures = [];
217 throw agg;
218 }
219 }
220
221 [Symbol.asyncDispose](): Promise<void> {
222 return this.all();
223 }
224}
225
226/** a promise that can be cancelled. */
227export type Cancelable<T> = Promise<T> & Disposable & {
228 cancel: (reason?: unknown) => void;
229};
230
231/** make an existing promise cancelable. */
232export function makeCancelable<T>(
233 promise: Promise<T>,
234 cancel: (reason?: unknown) => void,
235): Cancelable<T> {
236 const p = promise as Cancelable<T>;
237 p[Symbol.dispose] = cancel;
238 p.cancel = cancel;
239 return p;
240}
241
242/** wait `ms` milliseconds, then resolve. can be cancelled. */
243export function delay(ms: number): Cancelable<void> {
244 if (ms === Infinity) return makeCancelable(new Promise(() => {}), () => {});
245 let t: ts.Timer | null = null;
246 const maxTimer = 0x7FFFFFFF;
247 return makeCancelable(
248 new Promise(function tick(resolve) {
249 if (ms <= maxTimer) {
250 t = setTimeout(() => (t = null, resolve()), ms);
251 } else {
252 t = setTimeout(() => {
253 ms -= maxTimer;
254 tick(resolve);
255 }, maxTimer);
256 }
257 }),
258 () => {
259 if (t != null) clearTimeout(t);
260 },
261 );
262}
263
264export interface Deferred<T> extends AsyncDisposable, PromiseLike<T> {
265}
266
267/**
268 * Return something that can be awaited later without worry of an unhandled
269 * rejection. By disposing it, the promise exception will always be reported.
270 *
271 * ```ts
272 * // incorrect, possible unhandled rejection between a and b.
273 * const a = promise();
274 * const b = await promise();
275 * const c = await promise(b, await a);
276 * // correct
277 * await using a = async.deferred(promise());
278 * const b = await promise();
279 * const c = await promise(b, await a);
280 * ```
281 */
282export function deferred<T>(promise: Promise<T>): Deferred<T> {
283 let handled = false;
284 promise.catch(() => {}); // potential error is re-thrown on asyncDispose
285 return {
286 async [Symbol.asyncDispose]() {
287 if (!handled) await promise;
288 },
289 then(onFufilled, onRejected) {
290 return promise.finally(() => (handled = true)).then(onFufilled, onRejected);
291 },
292 };
293}
294
295import { ASSERT } from "./assert.ts";
296import * as ts from "./ts.ts";