Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
3.5 kB · 98 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899export interface ConsumerSchedulerOptions { concurrency: number; attempts?: number; retryDelayMs?: (attempt: number) => number; errorMessage?: string; onError?: (error: unknown, key: string, context: unknown) => void | Promise<void>;}
/** * Schedules independent logical consumers concurrently while preserving order * within each declaration/source key. */export class ConsumerScheduler { private readonly chains = new Map<string, Promise<void>>(); private readonly errors: unknown[] = []; private readonly withSlot: <T>(operation: () => Promise<T>) => Promise<T>; private readonly attempts: number; private readonly retryDelayMs: (attempt: number) => number; private readonly errorMessage: string; private readonly onError: ((error: unknown, key: string, context: unknown) => void | Promise<void>) | undefined;
constructor(options: ConsumerSchedulerOptions) { if (!Number.isSafeInteger(options.concurrency) || options.concurrency < 1) { throw new Error("Consumer scheduler concurrency must be a positive integer"); } this.attempts = options.attempts ?? 3; if (!Number.isSafeInteger(this.attempts) || this.attempts < 1) { throw new Error("Consumer scheduler attempts must be a positive integer"); } this.retryDelayMs = options.retryDelayMs ?? ((attempt) => 25 * attempt); this.errorMessage = options.errorMessage ?? "Consumer cycle failed"; this.onError = options.onError; this.withSlot = concurrentOperationLimiter(options.concurrency); }
enqueue(key: string, operation: () => Promise<void>, context?: unknown): void { const previous = this.chains.get(key) ?? Promise.resolve(); const current = previous.then(() => this.withSlot(() => this.runWithRetries(operation))).catch(async (error) => { this.errors.push(new Error(this.errorMessage)); if (this.onError) { try { await this.onError(error, key, context); } catch { this.errors.push(new Error("Consumer scheduler incident reporting failed")); } } }); this.chains.set(key, current); void current.then(() => { if (this.chains.get(key) === current) this.chains.delete(key); }); }
async drain(): Promise<void> { while (this.chains.size > 0) await Promise.all([...this.chains.values()]); if (this.errors.length > 0) throw new AggregateError(this.errors.splice(0), this.errorMessage); }
private async runWithRetries(operation: () => Promise<void>): Promise<void> { let lastError: unknown; for (let attempt = 1; attempt <= this.attempts; attempt += 1) { try { await operation(); return; } catch (error) { lastError = error; if (attempt < this.attempts) await delay(this.retryDelayMs(attempt)); } } throw lastError; }}
function concurrentOperationLimiter(maximum: number): <T>(operation: () => Promise<T>) => Promise<T> { let active = 0; const waiting: Array<() => void> = []; const acquire = async (): Promise<void> => { if (active < maximum) { active += 1; return; } await new Promise<void>((resolve) => waiting.push(resolve)); }; return async <T>(operation: () => Promise<T>): Promise<T> => { await acquire(); try { return await operation(); } finally { const next = waiting.shift(); if (next) next(); else active -= 1; } };}
async function delay(milliseconds: number): Promise<void> { await new Promise((resolve) => setTimeout(resolve, milliseconds));}