diff --git a/src/lib/scheduler.test.ts b/src/lib/scheduler.test.ts new file mode 100644 index 0000000..896016e --- /dev/null +++ b/src/lib/scheduler.test.ts @@ -0,0 +1,109 @@ +import { describe, expect, it } from 'vitest' +import { CoalescingQueue } from './scheduler' + +/** A task that resolves only when told to, recording that it started. */ +function gate(started: string[], name: string, value: T) { + let release!: (v: T) => void + const task = () => + new Promise((resolve) => { + started.push(name) + release = resolve + }) + return { task, release: () => release(value) } +} + +async function settled(): Promise { + // Let promise continuations run. + for (let i = 0; i < 5; i++) await Promise.resolve() +} + +describe('CoalescingQueue', () => { + it('never runs more tasks than the concurrency cap', async () => { + const q = new CoalescingQueue(2) + const started: string[] = [] + const g1 = gate(started, 't1', 'r1') + const g2 = gate(started, 't2', 'r2') + const g3 = gate(started, 't3', 'r3') + const results = [g1, g2, g3].map((g, i) => q.run(i, g.task)) + + await settled() + expect(started).toEqual(['t1', 't2']) + expect(q.depth).toBe(1) + + g1.release() + await settled() + expect(started).toEqual(['t1', 't2', 't3']) + + g2.release() + g3.release() + expect(await Promise.all(results)).toEqual(['r1', 'r2', 'r3']) + }) + + it('coalesces waiting tasks with the same key, newest task winning', async () => { + const q = new CoalescingQueue(1) + const started: string[] = [] + const blocker = gate(started, 'blocker', 'b') + void q.run(99, blocker.task) + + const first = q.run(7, async () => 'old') + const second = q.run(7, async () => 'new') + await settled() + expect(q.depth).toBe(1) // both key-7 requests share one queued entry + + blocker.release() + expect(await first).toBe('new') + expect(await second).toBe('new') + }) + + it('does not coalesce with a task that already started', async () => { + const q = new CoalescingQueue(1) + const started: string[] = [] + const running = gate(started, 'running', 'first') + const first = q.run(7, running.task) + await settled() + + const second = q.run(7, async () => 'second') + await settled() + expect(q.depth).toBe(1) + + running.release() + expect(await first).toBe('first') + expect(await second).toBe('second') + }) + + it('runs urgent tasks before earlier queued ones', async () => { + const q = new CoalescingQueue(1) + const started: string[] = [] + const blocker = gate(started, 'blocker', 'b') + void q.run(0, blocker.task) + + const slow = q.run(1, async () => { + started.push('slow') + return 'slow' + }) + const urgent = q.run(2, async () => { + started.push('urgent') + return 'urgent' + }, { urgent: true }) + + blocker.release() + await Promise.all([slow, urgent]) + expect(started).toEqual(['blocker', 'urgent', 'slow']) + }) + + it('rejects every caller of a coalesced entry when its task fails', async () => { + const q = new CoalescingQueue(1) + const started: string[] = [] + const blocker = gate(started, 'blocker', 'b') + void q.run(0, blocker.task) + + const first = q.run(7, async () => 'ok') + const second = q.run(7, async () => { + throw new Error('boom') + }) + blocker.release() + + await expect(first).rejects.toThrow('boom') + await expect(second).rejects.toThrow('boom') + }) +}) diff --git a/src/lib/scheduler.ts b/src/lib/scheduler.ts new file mode 100644 index 0000000..6bb00fb --- /dev/null +++ b/src/lib/scheduler.ts @@ -0,0 +1,68 @@ +// A small task queue for detection runs. Detection hits the network (well-known +// probes, record and subscription fetches), and some events fire it for many +// tabs at once — most notably the post-update content-script reinjection sweep, +// where every open tab reports hints within a second. The queue caps how many +// tasks run concurrently, coalesces queued tasks that share a key (a newer +// detection for a tab supersedes a queued older one), and lets user-facing +// requests jump ahead of background ones. + +interface Entry { + key: K + task: () => Promise + settlers: Array<{ resolve: (value: T) => void; reject: (reason: unknown) => void }> +} + +export class CoalescingQueue { + private readonly waiting: Array> = [] + private running = 0 + + constructor(private readonly concurrency: number) {} + + /** + * Enqueue `task`. If a task with the same key is still waiting, `task` + * replaces it and every caller awaiting that entry gets this task's result; + * a task that has already started is never disturbed. `urgent` tasks go to + * the front of the queue. + */ + run(key: K, task: () => Promise, opts: { urgent?: boolean } = {}): Promise { + return new Promise((resolve, reject) => { + let entry = this.waiting.find((e) => e.key === key) + if (entry) { + entry.task = task + entry.settlers.push({ resolve, reject }) + if (opts.urgent) { + this.waiting.splice(this.waiting.indexOf(entry), 1) + this.waiting.unshift(entry) + } + } else { + entry = { key, task, settlers: [{ resolve, reject }] } + if (opts.urgent) this.waiting.unshift(entry) + else this.waiting.push(entry) + } + this.pump() + }) + } + + /** Tasks queued but not yet started. */ + get depth(): number { + return this.waiting.length + } + + private pump(): void { + while (this.running < this.concurrency && this.waiting.length > 0) { + const entry = this.waiting.shift() + if (!entry) return + this.running++ + entry + .task() + .then( + (value) => entry.settlers.forEach((s) => s.resolve(value)), + (reason) => entry.settlers.forEach((s) => s.reject(reason)), + ) + .finally(() => { + this.running-- + this.pump() + }) + } + } +}