// Investor Flow — AdapterQueue (DESIGN.md §5). Shared queue for all sources. // Implements CacheScheduler (so CacheRepository.get schedules refreshes through it). // Dedupe: adapter_queue PRIMARY KEY(key) + status check collapses concurrent same-key jobs // to one fetch (ADR-0004). Token-bucket per source. Exponential backoff on failure. import type { DatabaseSync } from 'node:sqlite'; import type { CacheKey, SourceKind, CacheRepository, CacheScheduler } from '../cache/CacheRepository.ts'; import { parseCacheKey } from '../cache/CacheRepository.ts'; import type { SourceFetch } from '../adapters/SourceAdapter.ts'; export interface AdapterQueueOptions { db: DatabaseSync; adapters: Map; /** Override per-source min-interval (ms); tests pass 0 to skip throttling. */ rateLimitMs?: Partial>; } const DEFAULT_RATE_MS: Record = { yfinance: 1000, sec: 125, 'sec-fetch': 1000, reddit: 1000, x: 3000, macro: 1000, llm: 0 }; const BACKOFF_MS = [2000, 4000, 8000, 16000, 60000]; const MAX_ATTEMPTS = 5; const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); export interface QueueHealthDetailed { queued: number; in_flight: number; failed: number; backoff: number; done: number; backoff_until: string | null; last_error: string | undefined; paused: boolean; counts: Record; lastErrors: Array<{ key: string; error: string | null }>; } export interface QueueErrorEntry { attempt: number; errorMessage: string | null; errorStack: string | null; ts: string; } export interface ScheduleEntry { sourceKind: string; intervalMs: number; lastEnqueued: string | null; nextEnqueue: string | null; } export class AdapterQueue implements CacheScheduler { private readonly _db: DatabaseSync; private readonly _adapters: Map; private readonly _rate: Record; private _cache: CacheRepository | null = null; private _lastFetchAt: Partial> = {}; private _lastError: string | null = null; constructor(opts: AdapterQueueOptions) { this._db = opts.db; this._adapters = opts.adapters; this._rate = { ...DEFAULT_RATE_MS, ...(opts.rateLimitMs ?? {}) } as Record; } set cache(c: CacheRepository) { this._cache = c; } isPaused(): boolean { const row = this._db.prepare("SELECT value FROM queue_state WHERE key='paused'").get() as { value: string } | undefined; return row?.value === 'true'; } setPaused(paused: boolean): void { this._db.prepare("INSERT OR REPLACE INTO queue_state (key, value) VALUES ('paused', ?)").run(String(paused)); } async queue(key: CacheKey, scheduledFor?: string): Promise { const row = this._db.prepare('SELECT status, backoff_until FROM adapter_queue WHERE key=?').get(key) as { status?: string; backoff_until?: string | null } | undefined; if (row) { if (row.status === 'pending' || row.status === 'in_flight') return; if (row.status === 'backoff' && row.backoff_until && Date.parse(row.backoff_until) > Date.now()) return; } this._db.prepare('INSERT OR REPLACE INTO adapter_queue (key,status,last_attempt,retry_count,backoff_until,error,scheduled_for) VALUES (?,?,?,?,?,?,?)').run(key, 'pending', null, 0, null, null, scheduledFor ?? null); } async drain(): Promise { if (!this._cache) return; if (this.isPaused()) return; const now = Date.now(); const jobs = this._db.prepare( "SELECT key, retry_count, backoff_until, scheduled_for FROM adapter_queue WHERE status IN ('pending','backoff') ORDER BY (last_attempt IS NULL) DESC, last_attempt ASC", ).all() as Array<{ key: string; retry_count: number; backoff_until?: string | null; scheduled_for?: string | null }>; for (const job of jobs) { if (job.backoff_until && Date.parse(job.backoff_until) > now) continue; if (job.scheduled_for && Date.parse(job.scheduled_for) > now) continue; const { source } = parseCacheKey(job.key); const adapter = this._adapters.get(source); if (!adapter) { this._setStatus(job.key, 'done'); continue; } const last = this._lastFetchAt[source] ?? 0; const wait = (this._rate[source] ?? 0) - (Date.now() - last); if (wait > 0) await sleep(wait); this._lastFetchAt[source] = Date.now(); this._setStatus(job.key, 'in_flight'); const attempt = job.retry_count + 1; try { const res = await adapter.fetchOne(job.key); await this._cache.set(job.key, res.value, res.ttlClass, res.provenance); this._db.prepare("UPDATE adapter_queue SET status='done', last_attempt=?, error=NULL WHERE key=?").run(new Date().toISOString(), job.key); } catch (e) { const msg = e instanceof Error ? e.message : String(e); const stack = e instanceof Error ? (e.stack ?? msg) : msg; this._lastError = msg; this._db.prepare("INSERT INTO queue_errors (queue_key, attempt, error_message, error_stack, ts) VALUES (?,?,?,?,?)").run(job.key, attempt, msg, stack, new Date().toISOString()); if (attempt >= MAX_ATTEMPTS) { this._db.prepare("UPDATE adapter_queue SET status='failed', last_attempt=?, retry_count=?, backoff_until=NULL, error=? WHERE key=?").run(new Date().toISOString(), attempt, msg, job.key); } else { const bo = BACKOFF_MS[Math.min(attempt - 1, BACKOFF_MS.length - 1)]; this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=?, error=? WHERE key=?").run(new Date().toISOString(), attempt, new Date(Date.now() + bo).toISOString(), msg, job.key); } } } } retryJob(key: string): void { this._db.prepare("UPDATE adapter_queue SET status='pending', backoff_until=NULL, retry_count=0, error=NULL WHERE key=?").run(key); } retrySource(kind: string): number { const res = this._db.prepare( "UPDATE adapter_queue SET status='pending', backoff_until=NULL, retry_count=0, error=NULL WHERE status IN ('failed','backoff') AND substr(key,1,instr(key,':')-1)=?" ).run(kind); return Number(res.changes); } clearDone(olderThanMs: number = 86400000): number { const cutoff = new Date(Date.now() - olderThanMs).toISOString(); const res = this._db.prepare("DELETE FROM adapter_queue WHERE status='done' AND last_attempt < ?").run(cutoff); return Number(res.changes); } setSchedule(sourceKind: string, intervalMs: number): void { this._db.prepare("INSERT OR REPLACE INTO queue_schedules (source_kind, interval_ms, last_enqueued, next_enqueue) VALUES (?,?,?,?)").run(sourceKind, intervalMs, null, new Date(Date.now() + intervalMs).toISOString()); } deleteSchedule(sourceKind: string): void { this._db.prepare("DELETE FROM queue_schedules WHERE source_kind=?").run(sourceKind); } listSchedules(): ScheduleEntry[] { const rows = this._db.prepare("SELECT source_kind, interval_ms, last_enqueued, next_enqueue FROM queue_schedules ORDER BY source_kind").all() as Array<{ source_kind: string; interval_ms: number; last_enqueued: string | null; next_enqueue: string | null }>; return rows.map((r) => ({ sourceKind: r.source_kind, intervalMs: r.interval_ms, lastEnqueued: r.last_enqueued, nextEnqueue: r.next_enqueue, })); } async enqueueDueSchedules(): Promise { const now = new Date().toISOString(); const due = this._db.prepare("SELECT * FROM queue_schedules WHERE next_enqueue IS NOT NULL AND next_enqueue <= ?").all(now) as Array<{ source_kind: string; interval_ms: number; last_enqueued: string | null; next_enqueue: string | null }>; for (const s of due) { const symbols = this._db.prepare("SELECT symbol FROM symbol_demand WHERE in_demand=1").all() as Array<{ symbol: string }>; if (s.source_kind === 'sec-fetch') { for (const sym of symbols) { await this.queue(`sec-fetch:fetch:${sym.symbol}`); } } else if (s.source_kind === 'yfinance') { for (const sym of symbols) { await this.queue(`yfinance:quote:${sym.symbol}`); await this.queue(`yfinance:candles:${sym.symbol}:1d`); await this.queue(`yfinance:symbol:${sym.symbol}`); } } const nextEnqueue = new Date(Date.now() + s.interval_ms).toISOString(); this._db.prepare("UPDATE queue_schedules SET last_enqueued=?, next_enqueue=? WHERE source_kind=?").run(now, nextEnqueue, s.source_kind); } } seedDefaultSchedules(): void { const existing = this._db.prepare("SELECT COUNT(*) AS c FROM queue_schedules").get() as { c: number }; if (existing.c > 0) return; const defaults: Array<[string, number]> = [ ['sec-fetch', 86400000], ['yfinance', 300000], ]; const insert = this._db.prepare("INSERT OR IGNORE INTO queue_schedules (source_kind, interval_ms, last_enqueued, next_enqueue) VALUES (?,?,?,?)"); for (const [kind, ms] of defaults) { insert.run(kind, ms, null, new Date(Date.now() + ms).toISOString()); } } health(): QueueHealthDetailed { const paused = this.isPaused(); const row = this._db.prepare("SELECT SUM(status='pending') AS q, SUM(status='in_flight') AS i, SUM(status='failed') AS f, SUM(status='backoff') AS b, SUM(status='done') AS d, MAX(backoff_until) AS bu FROM adapter_queue").get() as { q: number | null; i: number | null; f: number | null; b: number | null; d: number | null; bu: string | null }; const failedJobs = this._db.prepare("SELECT key, error FROM adapter_queue WHERE status='failed' AND error IS NOT NULL ORDER BY last_attempt DESC LIMIT 5").all() as Array<{ key: string; error: string | null }>; return { queued: row.q ?? 0, in_flight: row.i ?? 0, failed: row.f ?? 0, backoff: row.b ?? 0, done: row.d ?? 0, backoff_until: row.bu ?? null, last_error: this._lastError ?? undefined, paused, counts: { pending: row.q ?? 0, in_flight: row.i ?? 0, failed: row.f ?? 0, backoff: row.b ?? 0, done: row.d ?? 0 }, lastErrors: failedJobs, }; } getErrorLog(key: string, limit: number = 10): QueueErrorEntry[] { const rows = this._db.prepare("SELECT attempt, error_message, error_stack, ts FROM queue_errors WHERE queue_key=? ORDER BY attempt DESC LIMIT ?").all(key, limit) as Array<{ attempt: number; error_message: string | null; error_stack: string | null; ts: string }>; return rows.map((r) => ({ attempt: r.attempt, errorMessage: r.error_message, errorStack: r.error_stack, ts: r.ts, })); } private _setStatus(key: string, status: string): void { this._db.prepare('UPDATE adapter_queue SET status=? WHERE key=?').run(status, key); } } export function createAdapterQueue(opts: AdapterQueueOptions): AdapterQueue { return new AdapterQueue(opts); }