Compare commits

..
Author SHA1 Message Date
Garry TanandClaude Fable 5 09f9b00f86 docs: restore parseNiceFlag docblock displaced by PRUNE_STATUSES; note --status/0d in top-level jobs help
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 11:20:16 -07:00
dfe1755da6 fix(jobs): prune --status filter + --older-than 0d as explicit no-age-floor (takeover of #2282)
Salvages the CLI plumbing from PR #2282 (queue.prune already accepted a
status[] param; the CLI exposed neither knob) and repairs the flaws found
in verification:

- --status completed,failed,dead,cancelled passes an explicit terminal
  subset through to queue.prune; anything else fails fast. Parsing lives
  in exported parsePruneStatuses (unit-tested, mirrors parseNiceFlag).
- --older-than 0d is documented and messaged as what it actually does:
  NO age floor — deletes ALL matching terminal jobs — not "same-day
  only" as the original PR body claimed. Help text + success line say so
  ("regardless of age") so an operator can't mistake it for a same-day
  cutoff.
- Real tests this time: queue-level status-filter + zero-age-floor cases
  in test/minions.test.ts, parser cases in test/jobs-prune-flags.test.ts
  (the original PR cited tests in a file that does not exist).

Co-authored-by: brettdavies <brettdavies@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:20:23 -07:00
11 changed files with 92 additions and 149 deletions
+1 -1
View File
@@ -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)
+1 -14
View File
@@ -1,5 +1,5 @@
import type { BrainEngine } from '../core/engine.ts';
import { embedBatch, currentEmbeddingSignature, resolveEmbeddingModelLabel } from '../core/embedding.ts';
import { embedBatch, currentEmbeddingSignature } from '../core/embedding.ts';
import type { ChunkInput } from '../core/types.ts';
import { chunkText } from '../core/chunkers/recursive.ts';
import { createProgress, type ProgressReporter } from '../core/progress.ts';
@@ -581,16 +581,11 @@ async function embedPage(
for (let j = 0; j < toEmbed.length; j++) {
embeddingMap.set(toEmbed[j].chunk_index, embeddings[j]);
}
// #1717: label each (re)embedded chunk with the model that actually
// produced its vector. Preserved chunks (not re-embedded this pass) keep
// their existing model so a mixed-model page isn't relabeled wholesale.
const embedModelLabel = resolveEmbeddingModelLabel();
const updated: ChunkInput[] = chunks.map(c => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: embeddingMap.get(c.chunk_index),
model: embeddingMap.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
@@ -722,16 +717,12 @@ async function embedAll(
for (let j = 0; j < toEmbed.length; j++) {
embeddingMap.set(toEmbed[j].chunk_index, embeddings[j]);
}
// #1717: stamp the resolved embedding model on (re)embedded chunks;
// preserve the existing model on chunks left untouched.
const embedModelLabel = resolveEmbeddingModelLabel();
// Preserve ALL chunks, only update embeddings for stale ones
const updated: ChunkInput[] = chunks.map(c => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: embeddingMap.get(c.chunk_index) ?? undefined,
model: embeddingMap.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, pageOpts));
@@ -1021,15 +1012,11 @@ async function embedAllStale(
for (let j = 0; j < stale.length; j++) {
staleIdxToEmbedding.set(stale[j].chunk_index, embeddings[j]);
}
// #1717: label the re-embedded (stale) chunks with the resolved
// model; preserve the existing model on the non-stale chunks.
const embedModelLabel = resolveEmbeddingModelLabel();
const merged: ChunkInput[] = existing.map(c => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
model: staleIdxToEmbedding.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
+33 -5
View File
@@ -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;
}
-7
View File
@@ -20,7 +20,6 @@
import type { BrainEngine } from './engine.ts';
import type { ChunkInput } from './types.ts';
import { embedBatchWithBackoff } from '../commands/embed.ts';
import { resolveEmbeddingModelLabel } from './embedding.ts';
import { type DbPacer, createNoopPacer, observed } from './db-pacer.ts';
import { AbortError } from './abort-check.ts';
@@ -201,17 +200,11 @@ export async function embedStaleForSource(
for (let j = 0; j < stale.length; j++) {
staleIdxToEmbedding.set(stale[j].chunk_index, embeddings[j]);
}
// #1717: label re-embedded chunks with the model that produced the
// vector; preserved chunks keep their existing model. Without this,
// upsertChunks falls back to DEFAULT_EMBEDDING_MODEL for every chunk
// (the same mislabel the embed.ts paths fixed).
const embedModelLabel = resolveEmbeddingModelLabel();
const merged: ChunkInput[] = existing.map((c) => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
model: staleIdxToEmbedding.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
// Carry through per-chunk metadata. upsertChunks writes these as
// EXCLUDED.<col> (not COALESCE), so omitting them here resets image
-15
View File
@@ -113,21 +113,6 @@ export async function embedBatch(
return results;
}
/**
* Resolve the embedding model label (`provider:model`) to stamp onto
* `content_chunks.model`, so each chunk records the model that actually
* produced its vector instead of the engine's hardcoded default (#1717).
* Returns undefined if the gateway is unconfigured; callers then fall back
* to the chunk's existing model rather than mislabeling it.
*/
export function resolveEmbeddingModelLabel(): string | undefined {
try {
return gatewayGetModel();
} catch {
return undefined;
}
}
/** Currently-configured embedding model (short form without provider prefix). */
export function getEmbeddingModelName(): string {
return gatewayGetModel().split(':').slice(1).join(':') || 'text-embedding-3-large';
+1 -11
View File
@@ -8,7 +8,7 @@ import { chunkText } from './chunkers/recursive.ts';
import { chunkCodeText, chunkCodeTextFull, detectCodeLanguage, CHUNKER_VERSION } from './chunkers/code.ts';
import { findChunkForOffset } from './chunkers/edge-extractor.ts';
import { extractCodeRefs, imageOfCandidates } from './link-extraction.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature, resolveEmbeddingModelLabel } from './embedding.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature } from './embedding.ts';
import { slugifyPath, slugifyCodePath, isCodeFilePath } from './sync.ts';
import type { ChunkInput, PageInput, PageType } from './types.ts';
import { computeEffectiveDate } from './effective-date.ts';
@@ -716,12 +716,8 @@ export async function importFromContent(
? chunks.map((c) => wrapChunkForEmbedding(c.chunk_text, prefix, c.chunk_source))
: chunks.map((c) => c.chunk_text);
const embeddings = await embedBatch(wrappedTexts);
// #1717: label each chunk with the model that actually produced its
// vector, not the engine's hardcoded default.
const embedModelLabel = resolveEmbeddingModelLabel();
for (let i = 0; i < chunks.length; i++) {
chunks[i].embedding = embeddings[i];
if (embedModelLabel) chunks[i].model = embedModelLabel;
// token_count tracks the wrapped string length so cost reporting
// reflects what we actually sent to the embedder.
chunks[i].token_count = Math.ceil(wrappedTexts[i].length / 4);
@@ -1145,10 +1141,7 @@ export async function importCodeFile(
const matched = existingByKey.get(key);
if (matched && matched.embedding) {
// Reuse the existing embedding verbatim. No API call, no cost.
// #1717: carry the existing model label along with the reused vector
// so the upsert doesn't relabel it with the engine default.
chunks[i]!.embedding = matched.embedding as Float32Array;
chunks[i]!.model = matched.model ?? undefined;
chunks[i]!.token_count = matched.token_count ?? undefined;
} else {
needsEmbedIndexes.push(i);
@@ -1160,12 +1153,9 @@ export async function importCodeFile(
try {
const textsToEmbed = needsEmbedIndexes.map((i) => chunks[i]!.chunk_text);
const embeddings = await embedBatch(textsToEmbed);
// #1717: stamp the model that produced these vectors.
const embedModelLabel = resolveEmbeddingModelLabel();
for (let j = 0; j < needsEmbedIndexes.length; j++) {
const i = needsEmbedIndexes[j]!;
chunks[i]!.embedding = embeddings[j]!;
if (embedModelLabel) chunks[i]!.model = embedModelLabel;
chunks[i]!.token_count = Math.ceil(chunks[i]!.chunk_text.length / 4);
}
} catch (e: unknown) {
-46
View File
@@ -15,7 +15,6 @@ import { describe, test, expect, beforeAll, afterAll, beforeEach } from 'bun:tes
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { resetPgliteState } from './helpers/reset-pglite.ts';
import { embedStaleForSource } from '../src/core/embed-stale.ts';
import { configureGateway, resetGateway } from '../src/core/ai/gateway.ts';
import type { ChunkInput } from '../src/core/types.ts';
let engine: PGLiteEngine;
@@ -277,49 +276,4 @@ describe('embedStaleForSource', () => {
// The stale text row actually got its embedding.
expect(txtRow.embedded_at).not.toBeNull();
});
// #1717: the backfill path must label re-embedded chunks with the model
// that produced the vector, and preserve the existing label on chunks it
// did not touch (before the fix, both were reset to the engine default).
test('labels re-embedded chunks with the gateway model, preserves untouched labels (#1717)', async () => {
configureGateway({
embedding_model: 'openai:text-embedding-3-large',
env: { OPENAI_API_KEY: 'sk-test-embed-stale-1717' },
});
try {
await engine.putPage('notes/model-label', {
type: 'note',
title: 'model-label',
compiled_truth: '# model-label\n\nseeded',
});
await engine.upsertChunks('notes/model-label', [
{
chunk_index: 0,
chunk_text: 'already embedded elsewhere',
chunk_source: 'compiled_truth',
embedding: new Float32Array(1536).fill(0.01),
model: 'voyage:voyage-3',
token_count: 4,
},
{
chunk_index: 1,
chunk_text: 'stale chunk needing embed',
chunk_source: 'compiled_truth',
token_count: 5,
embedding: undefined, // stale
},
]);
const result = await embedStaleForSource(engine, 'default', { embedFn: fakeEmbedFn });
expect(result.embedded).toBe(1);
const after = await engine.getChunks('notes/model-label');
const preserved = after.find((c) => c.chunk_index === 0)!;
const reembedded = after.find((c) => c.chunk_index === 1)!;
expect(reembedded.model).toBe('openai:text-embedding-3-large');
expect(preserved.model).toBe('voyage:voyage-3');
} finally {
resetGateway();
}
});
});
-33
View File
@@ -37,8 +37,6 @@ mock.module('../src/core/embedding.ts', () => ({
// setPageEmbeddingSignature / invalidateStaleSignatureEmbeddings resolve to
// null via the Proxy default, so the signature value is inert here.
currentEmbeddingSignature: () => 'test:model:1536',
// #1717: embed paths stamp this label on (re)embedded chunks.
resolveEmbeddingModelLabel: () => 'openai:text-embedding-3-large',
}));
// Import AFTER mocking.
@@ -805,34 +803,3 @@ describe('embedAllStale --source threading (D7)', () => {
expect((firstCallOpts as { sourceId?: string }).sourceId).toBe('media-corpus');
});
});
// #1717: content_chunks.model must record the model that actually produced
// each vector, not the gateway/engine default.
describe('content_chunks.model labeling (#1717)', () => {
test('stamps the resolved embedding model on re-embedded chunks, preserves it on untouched chunks', async () => {
let upserted: any[] | undefined;
// Chunk 0 is stale (no embedded_at) → gets re-embedded this pass.
// Chunk 1 is already embedded with a DIFFERENT model → must be preserved,
// not relabeled to the current model.
const chunks = [
{ chunk_index: 0, chunk_text: 'a', chunk_source: 'compiled_truth', embedded_at: null, model: 'zeroentropyai:zembed-1', token_count: 1 },
{ chunk_index: 1, chunk_text: 'b', chunk_source: 'compiled_truth', embedded_at: '2026-01-01', embedding: new Float32Array(1536), model: 'voyage:voyage-3', token_count: 1 },
];
const engine = mockEngine({
getPage: async () => ({ slug: 'notes/x', compiled_truth: 'a', timeline: '', source_id: 'default' }),
getChunks: async () => chunks,
upsertChunks: async (_slug: string, c: any[]) => { upserted = c; },
setPageEmbeddingSignature: async () => null,
});
await runEmbedCore(engine, { slugs: ['notes/x'] });
expect(upserted).toBeDefined();
const byIdx = Object.fromEntries(upserted!.map(c => [c.chunk_index, c]));
// Re-embedded chunk carries the model that produced its vector (was
// mislabeled with the default before the fix).
expect(byIdx[0].model).toBe('openai:text-embedding-3-large');
// Untouched chunk keeps its original model — no wholesale relabel.
expect(byIdx[1].model).toBe('voyage:voyage-3');
});
});
@@ -73,21 +73,4 @@ describe('importFromContent embedding_signature stamping (F1)', () => {
await importFromContent(engine, 'concepts/unstamped', '# Unstamped\n\nbody content.', { noEmbed: true });
expect(await signatureOf('concepts/unstamped')).toBeNull();
});
// #1717: content_chunks.model must record the model that produced the
// vector (the configured gateway model), not the engine's hardcoded
// default. The gateway here is configured to openai:text-embedding-3-large,
// which differs from DEFAULT_EMBEDDING_MODEL — so this fails without the
// import-path model stamping.
test('inline embed labels content_chunks.model with the configured model (#1717)', async () => {
await importFromContent(engine, 'concepts/labeled', '# Labeled\n\nsome body content to chunk and embed.', {});
const rows = await engine.executeRaw<{ model: string }>(
`SELECT cc.model FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = $1 AND p.source_id = 'default'`,
['concepts/labeled'],
);
expect(rows.length).toBeGreaterThan(0);
for (const r of rows) expect(r.model).toBe('openai:text-embedding-3-large');
});
});
+30
View File
@@ -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/);
});
});
+26
View File
@@ -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) ---