mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-16 01:42:23 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
89579780e0 | ||
|
|
cf2deedfc6 |
+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] [--status s,..] Clean old terminal jobs (0d = no age floor)
|
||||
jobs prune [--older-than 30d] Clean old jobs
|
||||
jobs stats Job health dashboard
|
||||
jobs work [--queue Q] Start worker daemon (Postgres only)
|
||||
|
||||
|
||||
+14
-1
@@ -1,5 +1,5 @@
|
||||
import type { BrainEngine } from '../core/engine.ts';
|
||||
import { embedBatch, currentEmbeddingSignature } from '../core/embedding.ts';
|
||||
import { embedBatch, currentEmbeddingSignature, resolveEmbeddingModelLabel } 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,11 +581,16 @@ 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),
|
||||
}));
|
||||
|
||||
@@ -717,12 +722,16 @@ 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));
|
||||
@@ -1012,11 +1021,15 @@ 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 }));
|
||||
|
||||
+5
-33
@@ -106,21 +106,6 @@ 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
|
||||
@@ -223,9 +208,7 @@ USAGE
|
||||
gbrain jobs get <id>
|
||||
gbrain jobs cancel <id>
|
||||
gbrain jobs retry <id>
|
||||
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 prune [--older-than 30d]
|
||||
gbrain jobs delete <id>
|
||||
gbrain jobs stats
|
||||
gbrain jobs smoke
|
||||
@@ -617,27 +600,16 @@ 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 non-negative number (days). Example: --older-than 30d; --older-than 0d removes the age floor (deletes ALL matching terminal jobs).');
|
||||
if (isNaN(days) || days <= 0) {
|
||||
console.error('Error: --older-than must be a positive number (days). Example: --older-than 30d');
|
||||
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),
|
||||
...(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}.`);
|
||||
const count = await queue.prune({ olderThan: new Date(Date.now() - days * 86400000) });
|
||||
console.log(`Pruned ${count} jobs older than ${days} days.`);
|
||||
break;
|
||||
}
|
||||
|
||||
|
||||
@@ -20,6 +20,7 @@
|
||||
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';
|
||||
|
||||
@@ -200,11 +201,17 @@ 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
|
||||
|
||||
@@ -113,6 +113,21 @@ 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';
|
||||
|
||||
+11
-1
@@ -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 } from './embedding.ts';
|
||||
import { embedBatch, embedMultimodal, currentEmbeddingSignature, resolveEmbeddingModelLabel } 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,8 +716,12 @@ 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);
|
||||
@@ -1141,7 +1145,10 @@ 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);
|
||||
@@ -1153,9 +1160,12 @@ 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) {
|
||||
|
||||
@@ -15,6 +15,7 @@ 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;
|
||||
@@ -276,4 +277,49 @@ 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();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
@@ -37,6 +37,8 @@ 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.
|
||||
@@ -803,3 +805,34 @@ 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,4 +73,21 @@ 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');
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,30 +0,0 @@
|
||||
/**
|
||||
* 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,32 +702,6 @@ 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