Files
investor-flow/app/server/src/queue/AdapterQueue.ts
T

100 lines
5.4 KiB
TypeScript
Raw Normal View History

// 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, AdapterHealth } from '../adapters/SourceAdapter.ts';
export interface AdapterQueueOptions {
db: DatabaseSync;
adapters: Map<SourceKind, SourceFetch>;
/** Override per-source min-interval (ms); tests pass 0 to skip throttling. */
rateLimitMs?: Partial<Record<SourceKind, number>>;
}
const DEFAULT_RATE_MS: Record<SourceKind, number> = { yfinance: 1000, sec: 125, reddit: 1000, x: 3000, macro: 1000, llm: 0 };
const BACKOFF_MS = [2000, 4000, 8000, 16000, 60000]; // yfinance policy §5 (2s..60s ceiling)
const MAX_ATTEMPTS = 5;
const sleep = (ms: number) => new Promise<void>((r) => setTimeout(r, ms));
export class AdapterQueue implements CacheScheduler {
private readonly _db: DatabaseSync;
private readonly _adapters: Map<SourceKind, SourceFetch>;
private readonly _rate: Record<SourceKind, number>;
private _cache: CacheRepository | null = null;
private _lastFetchAt: Partial<Record<SourceKind, number>> = {};
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<SourceKind, number>;
}
/** Link the cache (results are written here after a successful fetch). Breaks the cache<->queue cycle. */
set cache(c: CacheRepository) { this._cache = c; }
/** Enqueue a fetch. Dedupes against pending/in_flight; respects active backoff. force bypasses dedupe. */
async queue(key: CacheKey): Promise<void> {
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; // dedupe: already queued
if (row.status === 'backoff' && row.backoff_until && Date.parse(row.backoff_until) > Date.now()) return; // still backing off
}
this._db.prepare('INSERT OR REPLACE INTO adapter_queue (key,status,last_attempt,retry_count,backoff_until) VALUES (?,?,?,?,?)').run(key, 'pending', null, 0, null);
}
/** Process pending/backoff-expired jobs: fetch via the source adapter, write to cache, set done/backoff/failed. */
async drain(): Promise<void> {
if (!this._cache) return;
const now = Date.now();
const jobs = this._db.prepare(
"SELECT key, retry_count, backoff_until 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 }>;
for (const job of jobs) {
if (job.backoff_until && Date.parse(job.backoff_until) > now) continue; // backoff not yet expired
const { source } = parseCacheKey(job.key);
const adapter = this._adapters.get(source);
if (!adapter) { this._setStatus(job.key, 'done'); continue; } // no adapter registered -> don't leave stuck
// token bucket: enforce per-source min interval
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');
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=? WHERE key=?").run(new Date().toISOString(), job.key);
} catch (e) {
const retry = job.retry_count + 1;
this._lastError = e instanceof Error ? e.message : String(e);
if (retry >= MAX_ATTEMPTS) {
this._db.prepare("UPDATE adapter_queue SET status='failed', last_attempt=?, retry_count=?, backoff_until=NULL WHERE key=?").run(new Date().toISOString(), retry, job.key);
} else {
const bo = BACKOFF_MS[Math.min(retry - 1, BACKOFF_MS.length - 1)];
this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=? WHERE key=?").run(new Date().toISOString(), retry, new Date(Date.now() + bo).toISOString(), job.key);
}
}
}
}
health(): AdapterHealth {
const row = this._db.prepare("SELECT SUM(status='pending') AS q, SUM(status='in_flight') AS i, MAX(backoff_until) AS bu FROM adapter_queue").get() as { q: number | null; i: number | null; bu: string | null };
return { queued: row.q ?? 0, in_flight: row.i ?? 0, backoff_until: row.bu ?? null, last_error: this._lastError ?? undefined };
}
private _setStatus(key: string, status: string): void {
this._db.prepare('UPDATE adapter_queue SET status=? WHERE key=?').run(status, key);
}
}
/** Wire queue + cache together (breaks the cache<->scheduler cycle). Used by slice 1f. */
export function createAdapterQueue(opts: AdapterQueueOptions): AdapterQueue {
return new AdapterQueue(opts);
}