mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-17 10:22:34 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
09f9b00f86 | ||
|
|
dfe1755da6 |
+1
-1
@@ -2387,7 +2387,7 @@ JOBS (Minions)
|
||||
jobs get <id> Job details + history
|
||||
jobs cancel <id> Cancel job
|
||||
jobs retry <id> Re-queue failed/dead job
|
||||
jobs prune [--older-than 30d] Clean old jobs
|
||||
jobs prune [--older-than 30d] [--status s,..] Clean old terminal jobs (0d = no age floor)
|
||||
jobs stats Job health dashboard
|
||||
jobs work [--queue Q] Start worker daemon (Postgres only)
|
||||
|
||||
|
||||
+15
-33
@@ -107,14 +107,6 @@ export interface EmbedOpts {
|
||||
* runs lock every source in sorted order. dryRun skips it.
|
||||
*/
|
||||
singleFlight?: boolean;
|
||||
/**
|
||||
* #394: suppress human stdout summaries (the `[dry-run] Would embed ...` /
|
||||
* `Embedded N chunks ...` slog lines). Set by structured-output callers —
|
||||
* the cycle's embed phase (dream --json must keep stdout JSON-clean per
|
||||
* docs/progress-events.md) reports counts via its own PhaseResult instead.
|
||||
* Errors/warnings still go to stderr regardless.
|
||||
*/
|
||||
quiet?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -261,7 +253,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
|
||||
for (const s of opts.slugs) {
|
||||
if (isAborted(opts.signal)) break; // #1737: stop the per-slug loop on abort
|
||||
try {
|
||||
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal, opts.quiet);
|
||||
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal);
|
||||
} catch (e: unknown) {
|
||||
serr(` Error embedding ${s}: ${e instanceof Error ? e.message : e}`);
|
||||
}
|
||||
@@ -355,7 +347,6 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
|
||||
catchUp: opts.catchUp,
|
||||
pacer,
|
||||
paceMaxConcurrency,
|
||||
quiet: opts.quiet,
|
||||
}, opts.signal);
|
||||
} finally {
|
||||
// E1: surface pacing telemetry (human + structured) when pacing was on.
|
||||
@@ -385,7 +376,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
|
||||
return result;
|
||||
}
|
||||
if (opts.slug) {
|
||||
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal, opts.quiet);
|
||||
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal);
|
||||
return result;
|
||||
}
|
||||
throw new Error('No embed target specified. Pass { slug }, { slugs }, { all }, or { stale }.');
|
||||
@@ -530,7 +521,6 @@ async function embedPage(
|
||||
result: EmbedResult,
|
||||
sourceId?: string,
|
||||
signal?: AbortSignal,
|
||||
quiet?: boolean,
|
||||
) {
|
||||
const opts = sourceId ? { sourceId } : undefined;
|
||||
const page = await engine.getPage(slug, opts);
|
||||
@@ -575,7 +565,7 @@ async function embedPage(
|
||||
result.skipped += chunks.length - toEmbed.length;
|
||||
|
||||
if (toEmbed.length === 0) {
|
||||
if (!quiet) slog(`${slug}: all ${chunks.length} chunks already embedded`);
|
||||
slog(`${slug}: all ${chunks.length} chunks already embedded`);
|
||||
result.pages_processed++;
|
||||
return;
|
||||
}
|
||||
@@ -612,7 +602,7 @@ async function embedPage(
|
||||
}
|
||||
result.embedded += toEmbed.length;
|
||||
result.pages_processed++;
|
||||
if (!quiet) slog(`${slug}: embedded ${toEmbed.length} chunks`);
|
||||
slog(`${slug}: embedded ${toEmbed.length} chunks`);
|
||||
}
|
||||
|
||||
async function embedAll(
|
||||
@@ -630,8 +620,6 @@ async function embedAll(
|
||||
pacer?: DbPacer;
|
||||
/** Resolved concurrency cap (E-1: the worker count, no separate permit). */
|
||||
paceMaxConcurrency?: number;
|
||||
/** #394: suppress human stdout summaries (structured-output callers). */
|
||||
quiet?: boolean;
|
||||
},
|
||||
signal?: AbortSignal,
|
||||
) {
|
||||
@@ -775,12 +763,10 @@ async function embedAll(
|
||||
});
|
||||
|
||||
// Stdout summary preserved for scripts/tests that grep for counts.
|
||||
if (!staleOpts?.quiet) {
|
||||
if (dryRun) {
|
||||
slog(`[dry-run] Would embed ${result.would_embed} chunks across ${pages.length} pages`);
|
||||
} else {
|
||||
slog(`Embedded ${result.embedded} chunks across ${pages.length} pages`);
|
||||
}
|
||||
if (dryRun) {
|
||||
slog(`[dry-run] Would embed ${result.would_embed} chunks across ${pages.length} pages`);
|
||||
} else {
|
||||
slog(`Embedded ${result.embedded} chunks across ${pages.length} pages`);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -816,8 +802,6 @@ async function embedAllStale(
|
||||
pacer?: DbPacer;
|
||||
/** Resolved concurrency cap (E-1: the worker count, no separate permit). */
|
||||
paceMaxConcurrency?: number;
|
||||
/** #394: suppress human stdout summaries (structured-output callers). */
|
||||
quiet?: boolean;
|
||||
},
|
||||
signature?: string,
|
||||
externalSignal?: AbortSignal,
|
||||
@@ -835,7 +819,7 @@ async function embedAllStale(
|
||||
signature,
|
||||
...(sourceId && { sourceId }),
|
||||
});
|
||||
if (invalidated > 0 && !staleOpts?.quiet) {
|
||||
if (invalidated > 0) {
|
||||
slog(`[embed] invalidated ${invalidated} chunk(s) embedded under a prior model signature`);
|
||||
}
|
||||
}
|
||||
@@ -846,12 +830,10 @@ async function embedAllStale(
|
||||
dryRun && signature ? { ...sourceOpt, signature } : sourceOpt,
|
||||
);
|
||||
if (staleCount === 0) {
|
||||
if (!staleOpts?.quiet) {
|
||||
if (dryRun) {
|
||||
slog('[dry-run] Would embed 0 chunks (0 stale found)');
|
||||
} else {
|
||||
slog('Embedded 0 chunks (0 stale found)');
|
||||
}
|
||||
if (dryRun) {
|
||||
slog('[dry-run] Would embed 0 chunks (0 stale found)');
|
||||
} else {
|
||||
slog('Embedded 0 chunks (0 stale found)');
|
||||
}
|
||||
return;
|
||||
}
|
||||
@@ -860,7 +842,7 @@ async function embedAllStale(
|
||||
result.would_embed += staleCount;
|
||||
result.total_chunks += staleCount;
|
||||
if (onProgress) onProgress(1, 1, 0);
|
||||
if (!staleOpts?.quiet) slog(`[dry-run] Would embed ${staleCount} stale chunks`);
|
||||
slog(`[dry-run] Would embed ${staleCount} stale chunks`);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1100,7 +1082,7 @@ async function embedAllStale(
|
||||
if (budgetTimer) clearTimeout(budgetTimer);
|
||||
}
|
||||
|
||||
if (!staleOpts?.quiet) slog(`Embedded ${result.embedded} chunks across ${totalProcessedPages} pages`);
|
||||
slog(`Embedded ${result.embedded} chunks across ${totalProcessedPages} pages`);
|
||||
|
||||
// #1946 (OV2a): a catch-up pass that completed without being aborted but left
|
||||
// chunks unembedded means those chunks are stuck (a non-transient embed
|
||||
|
||||
+33
-5
@@ -106,6 +106,21 @@ export function parseMaxRssFlag(args: string[]): number | undefined {
|
||||
return parsed;
|
||||
}
|
||||
|
||||
/** Terminal statuses `jobs prune --status` accepts (PR #2282). Matches what
|
||||
* queue.prune can safely delete; anything else (waiting/active/…) is live. */
|
||||
export const PRUNE_STATUSES = ['completed', 'failed', 'dead', 'cancelled'] as const satisfies readonly MinionJobStatus[];
|
||||
|
||||
/** Parse a `--status a,b,c` value into prune statuses. Throws on any value
|
||||
* outside PRUNE_STATUSES (fail-fast, mirrors parseNiceValue). */
|
||||
export function parsePruneStatuses(raw: string): MinionJobStatus[] {
|
||||
const requested = raw.split(',').map(s => s.trim()).filter(Boolean);
|
||||
const invalid = requested.filter(s => !(PRUNE_STATUSES as readonly string[]).includes(s));
|
||||
if (requested.length === 0 || invalid.length > 0) {
|
||||
throw new Error(`--status accepts a comma-separated subset of [${PRUNE_STATUSES.join(', ')}]${invalid.length ? `. Invalid: ${invalid.join(', ')}` : ''}`);
|
||||
}
|
||||
return requested as MinionJobStatus[];
|
||||
}
|
||||
|
||||
/** Parse `--nice N` (then `GBRAIN_NICE` env). Returns:
|
||||
* - undefined if absent (no priority change — inherit)
|
||||
* - the validated integer in [-20, 19] otherwise
|
||||
@@ -208,7 +223,9 @@ USAGE
|
||||
gbrain jobs get <id>
|
||||
gbrain jobs cancel <id>
|
||||
gbrain jobs retry <id>
|
||||
gbrain jobs prune [--older-than 30d]
|
||||
gbrain jobs prune [--older-than 30d] [--status completed,failed,dead,cancelled]
|
||||
(--older-than 0d = no age floor: deletes ALL
|
||||
matching terminal jobs; pair with --status)
|
||||
gbrain jobs delete <id>
|
||||
gbrain jobs stats
|
||||
gbrain jobs smoke
|
||||
@@ -600,16 +617,27 @@ HANDLER TYPES (built in)
|
||||
case 'prune': {
|
||||
const olderThanStr = parseFlag(args, '--older-than') ?? '30d';
|
||||
const days = parseInt(olderThanStr, 10);
|
||||
if (isNaN(days) || days <= 0) {
|
||||
console.error('Error: --older-than must be a positive number (days). Example: --older-than 30d');
|
||||
if (isNaN(days) || days < 0) {
|
||||
console.error('Error: --older-than must be a non-negative number (days). Example: --older-than 30d; --older-than 0d removes the age floor (deletes ALL matching terminal jobs).');
|
||||
process.exit(1);
|
||||
}
|
||||
const statusFlag = parseFlag(args, '--status');
|
||||
let statuses: MinionJobStatus[] | undefined;
|
||||
if (statusFlag !== undefined) {
|
||||
try { statuses = parsePruneStatuses(statusFlag); }
|
||||
catch (e) { console.error(`Error: ${e instanceof Error ? e.message : String(e)}`); process.exit(1); }
|
||||
}
|
||||
|
||||
try { await queue.ensureSchema(); }
|
||||
catch (e) { console.error(e instanceof Error ? e.message : String(e)); process.exit(1); }
|
||||
|
||||
const count = await queue.prune({ olderThan: new Date(Date.now() - days * 86400000) });
|
||||
console.log(`Pruned ${count} jobs older than ${days} days.`);
|
||||
const count = await queue.prune({
|
||||
olderThan: new Date(Date.now() - days * 86400000),
|
||||
...(statuses ? { status: statuses } : {}),
|
||||
});
|
||||
const statusLabel = statuses ? statuses.join('+') : 'completed+dead+cancelled';
|
||||
const ageLabel = days === 0 ? 'regardless of age' : `older than ${days} days`;
|
||||
console.log(`Pruned ${count} ${statusLabel} jobs ${ageLabel}.`);
|
||||
break;
|
||||
}
|
||||
|
||||
|
||||
+1
-3
@@ -1214,9 +1214,7 @@ async function runPhaseEmbed(engine: BrainEngine, dryRun: boolean, signal?: Abor
|
||||
// 10-15 min one) bails within a batch instead of running to completion
|
||||
// after the job was killed — which left gbrain_cycle_locks held and
|
||||
// wedged every subsequent autopilot cycle.
|
||||
// #394: quiet — the cycle reports embed counts via its own PhaseResult;
|
||||
// raw `[dry-run] Would embed ...` stdout lines would corrupt `dream --json`.
|
||||
const result = await runEmbedCore(engine, { stale: true, dryRun, signal, quiet: true });
|
||||
const result = await runEmbedCore(engine, { stale: true, dryRun, signal });
|
||||
const embeddedCount = dryRun ? result.would_embed : result.embedded;
|
||||
return {
|
||||
phase: 'embed',
|
||||
|
||||
@@ -292,33 +292,6 @@ describe('runDream — output format', () => {
|
||||
expect(parsed).toHaveProperty('totals');
|
||||
});
|
||||
|
||||
// #394 / takeover of #854: the embed phase's `[dry-run] Would embed ...`
|
||||
// summary must not leak onto stdout ahead of the JSON CycleReport.
|
||||
test('--dry-run --json emits only JSON even when embed has stale chunks', async () => {
|
||||
await engine.putPage('concepts/testing', {
|
||||
type: 'concept',
|
||||
title: 'Testing',
|
||||
compiled_truth: 'Testing keeps JSON contracts honest.',
|
||||
timeline: '',
|
||||
});
|
||||
await engine.upsertChunks('concepts/testing', [
|
||||
{ chunk_index: 0, chunk_text: 'Testing keeps JSON contracts honest.', chunk_source: 'compiled_truth' },
|
||||
]);
|
||||
|
||||
const lines: string[] = [];
|
||||
const logSpy = spyOn(console, 'log').mockImplementation((msg: string) => { lines.push(String(msg)); });
|
||||
await runDream(engine, ['--dir', repo, '--phase', 'embed', '--dry-run', '--json']);
|
||||
logSpy.mockRestore();
|
||||
|
||||
const output = lines.join('\n');
|
||||
expect(output.trimStart().startsWith('{')).toBe(true);
|
||||
const parsed = JSON.parse(output);
|
||||
expect(parsed.schema_version).toBe('1');
|
||||
expect(parsed.phases[0].phase).toBe('embed');
|
||||
// The stale chunk was still counted in the structured report.
|
||||
expect(parsed.phases[0].details.would_embed).toBe(1);
|
||||
});
|
||||
|
||||
test('human output for clean status mentions "Brain is healthy"', async () => {
|
||||
const lines: string[] = [];
|
||||
const logSpy = spyOn(console, 'log').mockImplementation((msg: string) => { lines.push(String(msg)); });
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
/**
|
||||
* Unit tests for parsePruneStatuses (PR #2282) — `jobs prune --status` parsing.
|
||||
*/
|
||||
|
||||
import { describe, test, expect } from 'bun:test';
|
||||
import { parsePruneStatuses, PRUNE_STATUSES } from '../src/commands/jobs.ts';
|
||||
|
||||
describe('parsePruneStatuses', () => {
|
||||
test('parses a single status', () => {
|
||||
expect(parsePruneStatuses('failed')).toEqual(['failed']);
|
||||
});
|
||||
|
||||
test('parses a comma-separated list with whitespace', () => {
|
||||
expect(parsePruneStatuses(' completed, dead ')).toEqual(['completed', 'dead']);
|
||||
});
|
||||
|
||||
test('accepts every documented terminal status', () => {
|
||||
expect(parsePruneStatuses(PRUNE_STATUSES.join(','))).toEqual([...PRUNE_STATUSES]);
|
||||
});
|
||||
|
||||
test('throws on non-terminal statuses', () => {
|
||||
expect(() => parsePruneStatuses('waiting')).toThrow(/Invalid: waiting/);
|
||||
expect(() => parsePruneStatuses('completed,active')).toThrow(/Invalid: active/);
|
||||
});
|
||||
|
||||
test('throws on empty value', () => {
|
||||
expect(() => parsePruneStatuses('')).toThrow(/comma-separated subset/);
|
||||
expect(() => parsePruneStatuses(',')).toThrow(/comma-separated subset/);
|
||||
});
|
||||
});
|
||||
@@ -702,6 +702,32 @@ describe('MinionQueue: Prune', () => {
|
||||
const count = await queue.prune({ olderThan: new Date(Date.now() + 86400000) }); // future date = prune everything old enough
|
||||
expect(count).toBe(1); // only the cancelled one
|
||||
});
|
||||
|
||||
// PR #2282: `jobs prune --status` passes an explicit status subset through.
|
||||
test('status filter prunes only the requested terminal statuses', async () => {
|
||||
const cancelled = await queue.add('sync', {});
|
||||
await queue.cancelJob(cancelled.id);
|
||||
const dead = await queue.add('embed', {}, { max_attempts: 1 });
|
||||
await queue.claim('tok1', 30000, 'default', ['embed']);
|
||||
await queue.failJob(dead.id, 'tok1', 'boom', 'dead');
|
||||
|
||||
const count = await queue.prune({ olderThan: new Date(Date.now() + 86400000), status: ['dead'] });
|
||||
expect(count).toBe(1); // only the dead one
|
||||
|
||||
const remaining = await queue.getJobs({ status: 'cancelled' });
|
||||
expect(remaining.length).toBe(1);
|
||||
});
|
||||
|
||||
// PR #2282: `--older-than 0d` = no age floor — olderThan of "now" deletes
|
||||
// terminal jobs that finished moments ago.
|
||||
test('olderThan now (0d semantics) prunes just-terminated jobs', async () => {
|
||||
const job = await queue.add('sync', {});
|
||||
await queue.cancelJob(job.id);
|
||||
await new Promise(r => setTimeout(r, 5)); // ensure updated_at < now
|
||||
|
||||
const count = await queue.prune({ olderThan: new Date() });
|
||||
expect(count).toBe(1);
|
||||
});
|
||||
});
|
||||
|
||||
// --- Stats (1 test) ---
|
||||
|
||||
Reference in New Issue
Block a user