slice 1a-1d: DB + CacheRepository + YFinance adapter + AdapterQueue
Node 26 + node:sqlite backend (zero native deps; runtime glue adapted from Bun-spec design, deep-module architecture unchanged). 29 tests green across schema/CacheRepository/YFinance-parse/AdapterQueue-dedupe.
This commit is contained in:
@@ -0,0 +1,99 @@
|
||||
// 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);
|
||||
}
|
||||
Reference in New Issue
Block a user