feat: per-source pause/stop controls for queue admin
Add operator pause/resume/stop/start for individual adapter sources, persisted in queue_state. Pause freezes draining and retains jobs; stop clears pending/backoff jobs and blocks re-enqueue until start. Wires admin tRPC endpoints plus a clean two-row Actions layout in the queue admin page.
This commit is contained in:
@@ -72,6 +72,10 @@ export interface QueueHealthDetailed {
|
||||
lastErrors: Array<{ key: string; error: string | null }>;
|
||||
/** Active source-wide cool-downs (rate-limit pauses). */
|
||||
sourceCooldowns: SourceCooldownSnapshot[];
|
||||
/** Operator-paused sources (drain frozen, jobs retained). */
|
||||
pausedSources: string[];
|
||||
/** Operator-stopped sources (drain frozen, jobs cleared, enqueue blocked). */
|
||||
stoppedSources: string[];
|
||||
/** Pending jobs broken down by kind (quote/candles/symbol/…). */
|
||||
pendingByKind: Record<string, number>;
|
||||
demandSize: number;
|
||||
@@ -135,6 +139,83 @@ export class AdapterQueue implements CacheScheduler {
|
||||
this._db.prepare("INSERT OR REPLACE INTO queue_state (key, value) VALUES ('paused', ?)").run(String(paused));
|
||||
}
|
||||
|
||||
// ---- Per-source control (pause / stop) ---------------------------------
|
||||
// Persisted in queue_state like cool-downs: 'source_pause:<src>' = '1' and
|
||||
// 'source_stop:<src>' = '1'. Pause freezes draining but keeps queued jobs;
|
||||
// stop additionally clears pending/backoff jobs and blocks re-enqueue until
|
||||
// startSource() re-queues current demand.
|
||||
|
||||
private sourcePauseKey(source: SourceKind | string): string {
|
||||
return `source_pause:${source}`;
|
||||
}
|
||||
|
||||
private sourceStopKey(source: SourceKind | string): string {
|
||||
return `source_stop:${source}`;
|
||||
}
|
||||
|
||||
isSourcePaused(source: SourceKind | string): boolean {
|
||||
const row = this._db.prepare('SELECT value FROM queue_state WHERE key=?')
|
||||
.get(this.sourcePauseKey(source)) as { value: string } | undefined;
|
||||
return row?.value === '1';
|
||||
}
|
||||
|
||||
isSourceStopped(source: SourceKind | string): boolean {
|
||||
const row = this._db.prepare('SELECT value FROM queue_state WHERE key=?')
|
||||
.get(this.sourceStopKey(source)) as { value: string } | undefined;
|
||||
return row?.value === '1';
|
||||
}
|
||||
|
||||
/** Whether a source is paused OR stopped — both freeze draining. */
|
||||
isSourceControlled(source: SourceKind | string): boolean {
|
||||
return this.isSourcePaused(source) || this.isSourceStopped(source);
|
||||
}
|
||||
|
||||
pauseSource(source: SourceKind | string): void {
|
||||
this._db.prepare('INSERT OR REPLACE INTO queue_state (key, value) VALUES (?, ?)')
|
||||
.run(this.sourcePauseKey(source), '1');
|
||||
}
|
||||
|
||||
resumeSource(source: SourceKind | string): void {
|
||||
this._db.prepare("DELETE FROM queue_state WHERE key=?").run(this.sourcePauseKey(source));
|
||||
}
|
||||
|
||||
/** Stop a source: freeze draining, clear its pending/backoff jobs, block re-enqueue. */
|
||||
stopSource(source: SourceKind | string): number {
|
||||
this._db.prepare('INSERT OR REPLACE INTO queue_state (key, value) VALUES (?, ?)')
|
||||
.run(this.sourceStopKey(source), '1');
|
||||
this._db.prepare("DELETE FROM queue_state WHERE key=?").run(this.sourcePauseKey(source));
|
||||
const res = this._db.prepare(
|
||||
"DELETE FROM adapter_queue WHERE status IN ('pending','backoff') AND substr(key,1,instr(key,':')-1)=?",
|
||||
).run(String(source));
|
||||
return Number(res.changes);
|
||||
}
|
||||
|
||||
/** Start a source: clear pause+stop, then re-queue current demand via its schedules. */
|
||||
startSource(source: SourceKind | string): void {
|
||||
this._db.prepare("DELETE FROM queue_state WHERE key=?").run(this.sourceStopKey(source));
|
||||
this._db.prepare("DELETE FROM queue_state WHERE key=?").run(this.sourcePauseKey(source));
|
||||
// Force matching schedules to fire on the next enqueueDueSchedules tick.
|
||||
this._db.prepare("UPDATE queue_schedules SET next_enqueue=? WHERE source_kind=? OR source_kind LIKE ? OR source_kind LIKE ?")
|
||||
.run(new Date().toISOString(), String(source), `${source}-%`, `${source}:%`);
|
||||
}
|
||||
|
||||
listControlledSources(): { source: string; paused: boolean; stopped: boolean }[] {
|
||||
const rows = this._db.prepare(
|
||||
"SELECT key, value FROM queue_state WHERE key LIKE 'source_pause:%' OR key LIKE 'source_stop:%'",
|
||||
).all() as Array<{ key: string; value: string }>;
|
||||
const out = new Map<string, { source: string; paused: boolean; stopped: boolean }>();
|
||||
for (const r of rows) {
|
||||
const isStop = r.key.startsWith('source_stop:');
|
||||
const source = r.key.slice(isStop ? 'source_stop:'.length : 'source_pause:'.length);
|
||||
if (!source || r.value !== '1') continue;
|
||||
const cur = out.get(source) ?? { source, paused: false, stopped: false };
|
||||
if (isStop) cur.stopped = true;
|
||||
else cur.paused = true;
|
||||
out.set(source, cur);
|
||||
}
|
||||
return [...out.values()];
|
||||
}
|
||||
|
||||
/** Whether this source is under a rate-limit cool-down (no outbound calls). */
|
||||
isSourceCoolingDown(source: SourceKind | string, now = Date.now()): boolean {
|
||||
return this.getSourceCooldown(source, now).active;
|
||||
@@ -318,6 +399,13 @@ export class AdapterQueue implements CacheScheduler {
|
||||
const sym = id.split(':')[0]?.toUpperCase();
|
||||
if (sym && this.isSymbolQuarantined(sym)) return;
|
||||
|
||||
// Stopped sources accept no new jobs (schedules + manual triggers alike)
|
||||
// until startSource() clears the block and re-queues demand.
|
||||
try {
|
||||
const { source } = parseCacheKey(key);
|
||||
if (source && this.isSourceStopped(source)) return;
|
||||
} catch { /* malformed key — let the normal path decide */ }
|
||||
|
||||
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;
|
||||
@@ -451,13 +539,26 @@ export class AdapterQueue implements CacheScheduler {
|
||||
const used = kindUsed[kind] ?? 0;
|
||||
if (used >= kindBudget) continue;
|
||||
|
||||
// Prioritize x timeline jobs so tracked fund captures always make progress.
|
||||
// Without this, the yfinance backlog (tier 0-3) can starve x timeline jobs
|
||||
// indefinitely: they sort last and, under backlog pressure, the tier-3 skip
|
||||
// below drops them because their "sym" is a fund handle, not a symbol.
|
||||
const isXTimeline = source === 'x' && kind === 'timeline';
|
||||
|
||||
// Backlog pressure: defer non-critical yfinance kinds so quotes drain first.
|
||||
if (backlogPressure && source === 'yfinance' && NON_CRITICAL_YF_KINDS.has(kind)) continue;
|
||||
// Backlog pressure: skip T3 (background) symbols entirely so portfolio/alert symbols drain first.
|
||||
if (backlogPressure && sym && (tierMap.get(sym) ?? 99) >= 3) continue;
|
||||
// x timeline jobs are exempt: their sym is a fund handle (never in tierMap), so
|
||||
// without this they'd always be dropped while the yfinance backlog persists.
|
||||
if (backlogPressure && sym && (tierMap.get(sym) ?? 99) >= 3 && !isXTimeline) continue;
|
||||
|
||||
// Source-wide cool-down: skip all jobs for this vendor until the window ends.
|
||||
if (this.isSourceCoolingDown(source, Date.now())) continue;
|
||||
// Per-source operator control: paused/stopped sources never drain here.
|
||||
// (Stopped sources additionally have their pending/backoff rows deleted and
|
||||
// queue() blocks new enqueues, so this is belt-and-braces for in-flight
|
||||
// recoveries and anything queued just before a stop landed.)
|
||||
if (this.isSourceControlled(source)) continue;
|
||||
|
||||
if (family) {
|
||||
const famUsed = familyJobsThisDrain[family] ?? 0;
|
||||
@@ -483,6 +584,7 @@ export class AdapterQueue implements CacheScheduler {
|
||||
if (job.backoff_until && Date.parse(job.backoff_until) > now) continue;
|
||||
if (job.scheduled_for && Date.parse(job.scheduled_for) > now) continue;
|
||||
if (this.isSourceCoolingDown(source, Date.now())) continue;
|
||||
if (this.isSourceControlled(source)) continue;
|
||||
if (family && this.isVendorFamilyCoolingDown(family, Date.now())) continue;
|
||||
const kindBudget = DRAIN_KIND_BUDGET[kind] ?? DRAIN_KIND_BUDGET._default;
|
||||
if ((kindUsed[kind] ?? 0) >= kindBudget) continue;
|
||||
@@ -752,6 +854,15 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
continue;
|
||||
}
|
||||
|
||||
// Stopped sources: skip scheduling work entirely until startSource().
|
||||
// queue() already drops their jobs, so this just avoids wasted ticks and
|
||||
// the enqueue bookkeeping. Paused sources still enqueue (jobs sit queued).
|
||||
if (this.isSourceStopped(cooldownSource)) {
|
||||
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);
|
||||
continue;
|
||||
}
|
||||
|
||||
const symbols = this.demandSymbols();
|
||||
const d = this._db;
|
||||
|
||||
@@ -1116,11 +1227,16 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
} catch { /* ignore */ }
|
||||
|
||||
const cools = this.listSourceCooldowns();
|
||||
const controlled = this.listControlledSources();
|
||||
const notes: string[] = [];
|
||||
if (paused) notes.push('queue paused');
|
||||
for (const c of cools) {
|
||||
notes.push(`${c.source} cooling ${Math.ceil(c.remainingMs / 60_000)}m (hits=${c.consecutiveHits})`);
|
||||
}
|
||||
for (const c of controlled) {
|
||||
if (c.stopped) notes.push(`${c.source} stopped`);
|
||||
else if (c.paused) notes.push(`${c.source} paused`);
|
||||
}
|
||||
if ((row.q ?? 0) > 100) notes.push(`large backlog pending=${row.q}`);
|
||||
if (candleLagMs != null && candleLagMs > 3 * 86_400_000) notes.push(`SPY candle lag ${Math.round(candleLagMs / 86_400_000)}d`);
|
||||
if (demandSize > DEMAND_SET_SOFT_CAP) notes.push(`demand set ${demandSize} > soft cap ${DEMAND_SET_SOFT_CAP}`);
|
||||
@@ -1140,6 +1256,8 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
|
||||
counts: { pending: row.q ?? 0, in_flight: row.i ?? 0, failed: row.f ?? 0, backoff: row.b ?? 0, done: row.d ?? 0 },
|
||||
lastErrors: failedJobs,
|
||||
sourceCooldowns: cools,
|
||||
pausedSources: controlled.filter((c) => c.paused).map((c) => c.source),
|
||||
stoppedSources: controlled.filter((c) => c.stopped).map((c) => c.source),
|
||||
pendingByKind,
|
||||
demandSize,
|
||||
systemPinCount,
|
||||
|
||||
@@ -254,3 +254,69 @@ test('health() exposes pendingByKind and dataPlaneHealthy', async () => {
|
||||
assert.ok(h.pendingByKind.quote >= 1 || h.queued >= 1);
|
||||
assert.ok(Array.isArray(h.dataPlaneNotes));
|
||||
});
|
||||
|
||||
test('pauseSource freezes draining but keeps queued jobs; resume unfreezes', async () => {
|
||||
const { db, fake, queue } = setup();
|
||||
await queue.queue('yfinance:quote:NVDA');
|
||||
queue.pauseSource('yfinance');
|
||||
assert.equal(queue.isSourcePaused('yfinance'), true);
|
||||
assert.equal(queue.isSourceControlled('yfinance'), true);
|
||||
await queue.drain();
|
||||
assert.equal(statusOf(db, 'yfinance:quote:NVDA').status, 'pending', 'paused source must not drain');
|
||||
assert.equal(fake.calls.length, 0);
|
||||
queue.resumeSource('yfinance');
|
||||
assert.equal(queue.isSourcePaused('yfinance'), false);
|
||||
await queue.drain();
|
||||
assert.equal(statusOf(db, 'yfinance:quote:NVDA').status, 'done');
|
||||
assert.equal(fake.calls.length, 1);
|
||||
});
|
||||
|
||||
test('stopSource clears pending jobs and blocks new enqueues until start', async () => {
|
||||
const { db, fake, queue } = setup();
|
||||
await queue.queue('yfinance:quote:NVDA');
|
||||
await queue.queue('yfinance:quote:AAPL');
|
||||
const cleared = queue.stopSource('yfinance');
|
||||
assert.equal(cleared, 2, 'stop clears the pending rows');
|
||||
assert.equal(count(db, "status='pending'"), 0);
|
||||
assert.equal(queue.isSourceStopped('yfinance'), true);
|
||||
// queue() is a no-op for a stopped source
|
||||
await queue.queue('yfinance:quote:NVDA');
|
||||
assert.equal(count(db, "key='yfinance:quote:NVDA'"), 0, 'stopped source must reject new enqueues');
|
||||
// drain must not resurrect jobs (there are none), still no fetch
|
||||
await queue.drain();
|
||||
assert.equal(fake.calls.length, 0);
|
||||
// start re-enables the source
|
||||
queue.startSource('yfinance');
|
||||
assert.equal(queue.isSourceStopped('yfinance'), false);
|
||||
await queue.queue('yfinance:quote:NVDA');
|
||||
await queue.drain();
|
||||
assert.equal(statusOf(db, 'yfinance:quote:NVDA').status, 'done');
|
||||
});
|
||||
|
||||
test('stopped source schedules skip enqueue; startSource forces next_enqueue to now', async () => {
|
||||
const { db, queue } = setup();
|
||||
queue.seedDefaultSchedules();
|
||||
db.prepare("DELETE FROM adapter_queue").run();
|
||||
queue.stopSource('yfinance');
|
||||
// Force the yfinance-quote-watched schedule due — stopped source must skip enqueueing.
|
||||
db.prepare("UPDATE queue_schedules SET next_enqueue=? WHERE source_kind='yfinance-quote-watched'")
|
||||
.run(new Date(Date.now() - 1000).toISOString());
|
||||
await queue.enqueueDueSchedules();
|
||||
assert.equal(count(db, "status='pending'"), 0, 'stopped source schedules must not enqueue');
|
||||
queue.startSource('yfinance');
|
||||
const row = db.prepare("SELECT next_enqueue FROM queue_schedules WHERE source_kind='yfinance-quote-watched'").get() as { next_enqueue: string | null };
|
||||
assert.ok(row.next_enqueue, 'startSource must force the schedule to fire');
|
||||
assert.ok(Date.parse(row.next_enqueue) <= Date.now() + 5000);
|
||||
});
|
||||
|
||||
test('health() reports pausedSources and stoppedSources', async () => {
|
||||
const { queue } = setup();
|
||||
queue.pauseSource('fred');
|
||||
queue.stopSource('x');
|
||||
const h = queue.health();
|
||||
assert.ok(h.pausedSources.includes('fred'));
|
||||
assert.ok(h.stoppedSources.includes('x'));
|
||||
const notes = h.dataPlaneNotes.join(' ');
|
||||
assert.ok(notes.includes('fred paused'));
|
||||
assert.ok(notes.includes('x stopped'));
|
||||
});
|
||||
|
||||
@@ -1707,6 +1707,34 @@ const adminRouter = router({
|
||||
return { cleared };
|
||||
}),
|
||||
|
||||
queueSourcePause: adminProcedure
|
||||
.input(z.object({ source: z.string().min(1) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
ctx.queue.pauseSource(input.source);
|
||||
return { ok: true, source: input.source };
|
||||
}),
|
||||
|
||||
queueSourceResume: adminProcedure
|
||||
.input(z.object({ source: z.string().min(1) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
ctx.queue.resumeSource(input.source);
|
||||
return { ok: true, source: input.source };
|
||||
}),
|
||||
|
||||
queueSourceStop: adminProcedure
|
||||
.input(z.object({ source: z.string().min(1) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
const cleared = ctx.queue.stopSource(input.source);
|
||||
return { ok: true, source: input.source, cleared };
|
||||
}),
|
||||
|
||||
queueSourceStart: adminProcedure
|
||||
.input(z.object({ source: z.string().min(1) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
ctx.queue.startSource(input.source);
|
||||
return { ok: true, source: input.source };
|
||||
}),
|
||||
|
||||
queueClearDone: adminProcedure
|
||||
.input(z.object({ olderThanHours: z.number().min(1).max(336).optional().default(24) }))
|
||||
.mutation(({ ctx, input }) => {
|
||||
|
||||
Reference in New Issue
Block a user