mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-16 01:42:23 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
08f2397615 |
+2
-11
@@ -808,20 +808,12 @@ async function makeContext(engine: BrainEngine, params: Record<string, unknown>)
|
||||
// 'default'. Wrapped in try/catch so a doctor / single-source brain that
|
||||
// never set up sources still returns 'default' silently.
|
||||
let sourceId: string | undefined;
|
||||
// #2561: when the source resolved via a NON-explicit tier (path-match /
|
||||
// brain default / sole-non-default / seed default), unqualified search-shaped
|
||||
// reads span every `config.federated = true` source. Computed here (the
|
||||
// trusted local boundary) and consumed by federatedSearchScope in
|
||||
// operations.ts, which additionally gates on ctx.remote === false.
|
||||
let localFederated: string[] | undefined;
|
||||
try {
|
||||
const { resolveSourceWithTier, localFederatedSourceIds } = await import('./core/source-resolver.ts');
|
||||
const { resolveSourceId } = await import('./core/source-resolver.ts');
|
||||
// params.source is set when a CLI flag was parsed for the op (rare; most
|
||||
// CLI ops don't take --source). Falls through to env/dotfile/path-match.
|
||||
const explicit = (params.source as string | undefined) ?? null;
|
||||
const resolved = await resolveSourceWithTier(engine, explicit);
|
||||
sourceId = resolved.source_id;
|
||||
localFederated = await localFederatedSourceIds(engine, resolved.source_id, resolved.tier);
|
||||
sourceId = await resolveSourceId(engine, explicit);
|
||||
} catch {
|
||||
// Source resolution failed (e.g. sources table doesn't exist on a fresh
|
||||
// pre-init brain). Leave sourceId unset; engine read methods fall through
|
||||
@@ -842,7 +834,6 @@ async function makeContext(engine: BrainEngine, params: Record<string, unknown>)
|
||||
// table). Matches dispatch.ts's auto-fill so the contract holds across
|
||||
// every transport.
|
||||
sourceId: sourceId ?? 'default',
|
||||
...(localFederated ? { localFederatedSourceIds: localFederated } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -1651,7 +1651,7 @@ async function extractTimelineFromDB(
|
||||
* make re-extraction idempotent). EVERY processed page is stamped, including
|
||||
* zero-link pages — they WERE processed.
|
||||
*/
|
||||
async function extractStaleFromDB(
|
||||
export async function extractStaleFromDB(
|
||||
engine: BrainEngine,
|
||||
opts: {
|
||||
dryRun: boolean;
|
||||
|
||||
+39
-1
@@ -1479,7 +1479,31 @@ export async function registerBuiltinHandlers(
|
||||
embedSkipReason = 'auto_embed_disabled';
|
||||
}
|
||||
|
||||
return { ...result, embed_job_id: embedJobId, embed_skip_reason: embedSkipReason };
|
||||
// #2849: large-sync extract deferral follow-up. performSync skips inline
|
||||
// link/timeline extraction when totalChanges > 100, leaving
|
||||
// links_extracted_at unstamped. A standalone sync job (webhook push,
|
||||
// sync trigger) has no autopilot extract phase behind it, so the pages
|
||||
// would stay extraction-stale until a manual `gbrain extract --stale`.
|
||||
// Queue a source-scoped stale sweep instead. Best-effort + idempotent:
|
||||
// a duplicate sweep finds 0 stale pages and no-ops.
|
||||
let extractJobId: number | null = null;
|
||||
if (result.extractDeferred) {
|
||||
try {
|
||||
const { MinionQueue } = await import('../core/minions/queue.ts');
|
||||
const queue = new MinionQueue(engine);
|
||||
const followUp = await queue.add(
|
||||
'extract',
|
||||
{ stale: true, ...(sourceId ? { sourceId } : {}) },
|
||||
{
|
||||
idempotency_key: `sync-extract-stale:${sourceId ?? 'default'}:${Math.floor(Date.now() / 30_000)}`,
|
||||
maxWaiting: 1,
|
||||
},
|
||||
);
|
||||
extractJobId = followUp.id;
|
||||
} catch { /* best-effort: extract --stale sweeps it later */ }
|
||||
}
|
||||
|
||||
return { ...result, embed_job_id: embedJobId, embed_skip_reason: embedSkipReason, extract_stale_job_id: extractJobId };
|
||||
});
|
||||
|
||||
registerBuiltinJob(worker, engine, 'embed', async (job) => {
|
||||
@@ -1652,6 +1676,20 @@ export async function registerBuiltinHandlers(
|
||||
});
|
||||
|
||||
worker.register('extract', async (job) => {
|
||||
// #2849: stale-sweep mode — the sync handler's large-sync deferral
|
||||
// follow-up. DB-source (reads page content from the DB, so it runs on
|
||||
// checkout-less brains), source-scopable, idempotent. Same core as
|
||||
// `gbrain extract --stale`.
|
||||
if (job.data.stale === true) {
|
||||
const { extractStaleFromDB } = await import('./extract.ts');
|
||||
return await extractStaleFromDB(engine, {
|
||||
dryRun: !!job.data.dryRun,
|
||||
jsonMode: false,
|
||||
includeFrontmatter: false,
|
||||
sourceIdFilter: typeof job.data.sourceId === 'string' ? job.data.sourceId : undefined,
|
||||
catchUp: false,
|
||||
});
|
||||
}
|
||||
const { runExtractCore } = await import('./extract.ts');
|
||||
const mode = (typeof job.data.mode === 'string' && ['links', 'timeline', 'all'].includes(job.data.mode))
|
||||
? (job.data.mode as 'links' | 'timeline' | 'all')
|
||||
|
||||
@@ -2146,8 +2146,13 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
// Other event types (ping, pull_request, etc.) return 202 'ignored'
|
||||
// so GitHub doesn't retry.
|
||||
// D15.5: HMAC compare uses the shared safeHexEqual helper.
|
||||
// D18: submits 'sync' job with auto_embed_backfill=true and priority -10
|
||||
// (above autopilot's 0).
|
||||
// D18: submits 'sync' job with extraction + auto_embed_backfill enabled and
|
||||
// priority -10 (above autopilot's 0). noExtract:false opts normal
|
||||
// incremental pushes into sync's inline link/timeline extraction (#2849
|
||||
// — the standalone sync handler defaults noExtract to TRUE, which left
|
||||
// webhook-imported pages permanently stale). Large (>100 file) pushes
|
||||
// defer inline extract; the sync handler queues an extract --stale
|
||||
// follow-up job for that branch.
|
||||
// ---------------------------------------------------------------------------
|
||||
const githubWebhookLimiter = rateLimit({
|
||||
windowMs: 60_000,
|
||||
@@ -2267,6 +2272,7 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
'sync',
|
||||
{
|
||||
sourceId: source.id,
|
||||
noExtract: false,
|
||||
auto_embed_backfill: true,
|
||||
embed_reason: 'webhook',
|
||||
},
|
||||
|
||||
+19
-1
@@ -222,6 +222,14 @@ export interface SyncResult {
|
||||
* everything," the exact misdiagnosis in the #1794 recurrence report.
|
||||
*/
|
||||
bankedFiles?: number;
|
||||
/**
|
||||
* #2849: true when extraction was REQUESTED (noExtract false) but this sync
|
||||
* skipped inline link/timeline extraction because totalChanges > 100 (the
|
||||
* #1794 large-sync deferral). links_extracted_at stays unstamped for the
|
||||
* imported pages. The standalone `sync` job handler queues a source-scoped
|
||||
* `extract --stale` follow-up when set; CLI runs print the manual hint.
|
||||
*/
|
||||
extractDeferred?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1379,6 +1387,10 @@ See also:
|
||||
{
|
||||
sourceId: sourceIdArg,
|
||||
repoPath: source.local_path,
|
||||
// #2849: opt in to inline extraction — the standalone sync handler
|
||||
// defaults noExtract to TRUE (dedupe for doctor's [sync, extract]
|
||||
// remediation plan), which would leave triggered syncs extraction-stale.
|
||||
noExtract: false,
|
||||
auto_embed_backfill: true,
|
||||
embed_reason: 'sync_trigger',
|
||||
},
|
||||
@@ -3287,11 +3299,16 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
|
||||
// the stale sweep scans the whole source, so banked-across-runs pages are
|
||||
// covered regardless.
|
||||
const extractOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
|
||||
let extractDeferred = false;
|
||||
if (!opts.noExtract && totalChanges > 100 && pagesAffected.length > 0) {
|
||||
// #2849: surface the deferral to callers. A standalone sync job (webhook
|
||||
// push, sync trigger) has no autopilot extract phase behind it, so the
|
||||
// job handler queues an `extract --stale` follow-up off this flag.
|
||||
extractDeferred = true;
|
||||
slog(
|
||||
` Large sync: deferring link/timeline extraction. ` +
|
||||
`Run 'gbrain extract --stale${opts.sourceId ? ` --source-id ${opts.sourceId}` : ''}' ` +
|
||||
`(or let the autopilot cycle's extract phase sweep it).`,
|
||||
`(sync jobs queue this follow-up automatically).`,
|
||||
);
|
||||
}
|
||||
if (!opts.noExtract && totalChanges <= 100 && pagesAffected.length > 0) {
|
||||
@@ -3400,6 +3417,7 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
|
||||
chunksCreated,
|
||||
embedded,
|
||||
pagesAffected,
|
||||
extractDeferred,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
+2
-61
@@ -424,23 +424,6 @@ export interface OperationContext {
|
||||
* satisfied even on single-source brains.
|
||||
*/
|
||||
sourceId: string;
|
||||
/**
|
||||
* #2561 — federated read scope for UNQUALIFIED local CLI reads.
|
||||
*
|
||||
* Set ONLY by the local CLI's context builder (src/cli.ts makeContext), and
|
||||
* only when the source resolved via a non-explicit tier (local_path /
|
||||
* brain_default / sole_non_default / seed_default — NOT --source, NOT
|
||||
* GBRAIN_SOURCE, NOT a .gbrain-source dotfile). Contains the resolved
|
||||
* source first, then every other `config.federated = true` source, so an
|
||||
* unqualified `gbrain search "X"` spans federated sources as
|
||||
* docs/guides/multi-source-brains.md promises.
|
||||
*
|
||||
* Consumed exclusively by `federatedSearchScope` and ONLY when
|
||||
* `ctx.remote === false` — a remote caller's scope stays governed by
|
||||
* `ctx.auth.allowedSources` / scalar `ctx.sourceId` (source-isolation
|
||||
* invariant, fail-closed).
|
||||
*/
|
||||
localFederatedSourceIds?: string[];
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -556,45 +539,6 @@ export function resolveRequestedScope(
|
||||
return sourceScopeOpts(ctx);
|
||||
}
|
||||
|
||||
/**
|
||||
* #2561 — source scope for the search-shaped read ops (`search`, `query`).
|
||||
*
|
||||
* Delegates to `resolveRequestedScope` (the single trust+grant resolver), then
|
||||
* widens an UNQUALIFIED trusted-local scalar scope to the CLI-computed
|
||||
* federated set (`ctx.localFederatedSourceIds`, resolved source first). This is
|
||||
* what makes `sources add --federated` mean something for local search: a
|
||||
* federated source participates in unqualified `gbrain search "X"` results.
|
||||
*
|
||||
* The expansion NEVER applies when:
|
||||
* - the caller is not strictly trusted-local (`ctx.remote !== false`) —
|
||||
* remote scope stays grant-governed (fail-closed source isolation);
|
||||
* - a per-call `source_id` was passed (explicit wins, including `__all__`);
|
||||
* - the resolver already produced a federated array (OAuth grant);
|
||||
* - the CLI resolved the source from an explicit signal (--source / env /
|
||||
* dotfile) — makeContext leaves `localFederatedSourceIds` unset then.
|
||||
*
|
||||
* Deliberately NOT inside `sourceScopeOpts`: code-intel ops collapse a
|
||||
* multi-element scope to an error (`resolveCodeIntelScope`), and non-search
|
||||
* reads (get_page, get_links, …) keep their long-standing scalar behavior.
|
||||
*/
|
||||
export function federatedSearchScope(
|
||||
ctx: OperationContext,
|
||||
sourceIdParam?: string,
|
||||
): { sourceId?: string; sourceIds?: string[] } {
|
||||
const scope = resolveRequestedScope(ctx, sourceIdParam);
|
||||
if (
|
||||
ctx.remote === false &&
|
||||
sourceIdParam === undefined &&
|
||||
scope.sourceId !== undefined &&
|
||||
scope.sourceIds === undefined &&
|
||||
ctx.localFederatedSourceIds !== undefined &&
|
||||
ctx.localFederatedSourceIds.length > 1
|
||||
) {
|
||||
return { sourceIds: ctx.localFederatedSourceIds };
|
||||
}
|
||||
return scope;
|
||||
}
|
||||
|
||||
/**
|
||||
* Code-intel adapter for `resolveRequestedScope`. Graph traversal
|
||||
* (code_callers/code_callees/code_blast/code_flow) is single-source by design —
|
||||
@@ -1504,8 +1448,7 @@ const search: Operation = {
|
||||
const queryText = p.query as string;
|
||||
const limit = (p.limit as number) || 20;
|
||||
const offset = (p.offset as number) || 0;
|
||||
// #2561: unqualified trusted-local search spans federated sources.
|
||||
const scope = federatedSearchScope(ctx);
|
||||
const scope = sourceScopeOpts(ctx);
|
||||
|
||||
// T4/D5 — per-call mode honored ONLY for trusted/local callers so a remote
|
||||
// OAuth client can't escalate to the costly tokenmax bundle. Local + unknown
|
||||
@@ -1667,9 +1610,7 @@ const query: Operation = {
|
||||
// is spread into BOTH the image-similarity searchVector path and the text
|
||||
// hybridSearch path below, so both honor the same grant.
|
||||
const sourceIdParam = typeof p.source_id === 'string' ? p.source_id : undefined;
|
||||
// #2561: unqualified trusted-local query spans federated sources (per-call
|
||||
// source_id / remote grants still resolve through resolveRequestedScope).
|
||||
const querySourceScope = federatedSearchScope(ctx, sourceIdParam);
|
||||
const querySourceScope = resolveRequestedScope(ctx, sourceIdParam);
|
||||
|
||||
// v0.27.1: image-similarity branch. Bypasses hybridSearch (which is
|
||||
// text-only); embeds the image via embedMultimodal and runs a direct
|
||||
|
||||
@@ -353,45 +353,6 @@ export async function resolveSourceWithTier(
|
||||
return { source_id: 'default', tier: 'seed_default' };
|
||||
}
|
||||
|
||||
/**
|
||||
* #2561 — compute the federated read scope for an UNQUALIFIED local CLI call.
|
||||
*
|
||||
* `sources add --federated` promises that a `config.federated = true` source
|
||||
* "participates in unqualified `gbrain search` results"
|
||||
* (docs/guides/multi-source-brains.md). This helper turns that promise into a
|
||||
* scope: given the resolved source and WHICH tier resolved it, return
|
||||
* `[resolvedSource, ...other federated source ids]` — or `undefined` when the
|
||||
* expansion must not apply:
|
||||
*
|
||||
* - explicit tiers (`flag` / `env` / `dotfile`): the user named a source;
|
||||
* scalar scope stands (that IS the qualified case);
|
||||
* - no other federated source exists: keep the scalar fast path unchanged.
|
||||
*
|
||||
* Archived sources are excluded (same rationale as pickSoleNonDefaultSource);
|
||||
* the archived column is v34+, so fall back to the un-archived query on older
|
||||
* brains. Callers put the result on `OperationContext.localFederatedSourceIds`
|
||||
* — consumed only by `federatedSearchScope` and only when `remote === false`.
|
||||
*/
|
||||
export async function localFederatedSourceIds(
|
||||
engine: BrainEngine,
|
||||
sourceId: string,
|
||||
tier: SourceTier,
|
||||
): Promise<string[] | undefined> {
|
||||
if (tier === 'flag' || tier === 'env' || tier === 'dotfile') return undefined;
|
||||
let rows: Array<{ id: string }>;
|
||||
try {
|
||||
rows = await engine.executeRaw<{ id: string }>(
|
||||
`SELECT id FROM sources WHERE config->>'federated' = 'true' AND archived = false ORDER BY id`,
|
||||
);
|
||||
} catch {
|
||||
rows = await engine.executeRaw<{ id: string }>(
|
||||
`SELECT id FROM sources WHERE config->>'federated' = 'true' ORDER BY id`,
|
||||
);
|
||||
}
|
||||
const ids = [sourceId, ...rows.map((r) => r.id).filter((id) => id !== sourceId)];
|
||||
return ids.length > 1 ? ids : undefined;
|
||||
}
|
||||
|
||||
/** Exposed for tests. */
|
||||
export const __testing = {
|
||||
readDotfileWalk,
|
||||
|
||||
@@ -1,154 +0,0 @@
|
||||
/**
|
||||
* #2561 — sources.config.federated participates in UNQUALIFIED local CLI
|
||||
* search/query.
|
||||
*
|
||||
* Pre-fix: the local CLI always emitted a scalar `{sourceId}` scope (required
|
||||
* field, auto-filled 'default'), so a source registered with
|
||||
* `gbrain sources add --federated` was invisible to an unqualified
|
||||
* `gbrain search "X"` — contradicting docs/guides/multi-source-brains.md
|
||||
* ("Source participates in unqualified `gbrain search` results").
|
||||
*
|
||||
* Fix: the CLI context builder computes `ctx.localFederatedSourceIds`
|
||||
* (resolved source + every other federated source) whenever the source
|
||||
* resolved via a NON-explicit tier; `federatedSearchScope` widens the scalar
|
||||
* scope to that set for the `search` / `query` ops — trusted-local only
|
||||
* (`ctx.remote === false`), never for remote callers, never when a per-call
|
||||
* `source_id` or an explicit --source/env/dotfile was given.
|
||||
*/
|
||||
import { describe, test, expect, beforeAll, afterAll } from 'bun:test';
|
||||
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
|
||||
import { localFederatedSourceIds } from '../src/core/source-resolver.ts';
|
||||
import {
|
||||
federatedSearchScope,
|
||||
operations,
|
||||
type OperationContext,
|
||||
} from '../src/core/operations.ts';
|
||||
|
||||
let engine: PGLiteEngine;
|
||||
const search = operations.find((o) => o.name === 'search')!;
|
||||
|
||||
function ctxOf(overrides: Partial<OperationContext> = {}): OperationContext {
|
||||
return {
|
||||
engine: engine as any,
|
||||
config: {} as any,
|
||||
logger: console as any,
|
||||
dryRun: false,
|
||||
remote: false,
|
||||
sourceId: 'default',
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
beforeAll(async () => {
|
||||
engine = new PGLiteEngine();
|
||||
await engine.connect({});
|
||||
await engine.initSchema();
|
||||
// Seeded 'default' source is federated=true. Add:
|
||||
// wiki — federated (must join unqualified search)
|
||||
// private — NOT federated (must stay invisible unless explicitly named)
|
||||
// oldnews — federated but archived (must stay excluded)
|
||||
await engine.executeRaw(
|
||||
`INSERT INTO sources (id, name, local_path, config) VALUES ('wiki', 'wiki', '/tmp/wiki', '{"federated": true}'::jsonb)`,
|
||||
);
|
||||
await engine.executeRaw(
|
||||
`INSERT INTO sources (id, name, local_path, config) VALUES ('private', 'private', '/tmp/private', '{}'::jsonb)`,
|
||||
);
|
||||
await engine.executeRaw(
|
||||
`INSERT INTO sources (id, name, local_path, config, archived) VALUES ('oldnews', 'oldnews', '/tmp/oldnews', '{"federated": true}'::jsonb, true)`,
|
||||
);
|
||||
const pages: Array<[slug: string, sourceId: string, where: string]> = [
|
||||
['notes/home', 'default', 'default'],
|
||||
['wiki/topic', 'wiki', 'wiki'],
|
||||
['private/topic', 'private', 'private'],
|
||||
['old/topic', 'oldnews', 'oldnews'],
|
||||
];
|
||||
for (const [slug, sourceId, where] of pages) {
|
||||
await engine.putPage(slug, {
|
||||
type: 'note', title: `Topic in ${where}`, compiled_truth: `the zebra telescope in ${where}`, frontmatter: {},
|
||||
}, { sourceId });
|
||||
await engine.upsertChunks(slug, [
|
||||
{ chunk_index: 0, chunk_text: `the zebra telescope in ${where}`, chunk_source: 'compiled_truth' },
|
||||
], { sourceId });
|
||||
}
|
||||
// Keyword-only search path: no embedding provider needed in tests.
|
||||
await engine.setConfig('search.mcp_keyword_only', 'true');
|
||||
}, 60_000);
|
||||
|
||||
afterAll(async () => {
|
||||
if (engine) await engine.disconnect();
|
||||
}, 60_000);
|
||||
|
||||
describe('localFederatedSourceIds — CLI-side scope computation', () => {
|
||||
test('non-explicit tier: resolved source first, then other federated, archived excluded', async () => {
|
||||
expect(await localFederatedSourceIds(engine, 'default', 'seed_default')).toEqual(['default', 'wiki']);
|
||||
});
|
||||
|
||||
test('non-federated resolved source still joins its own scope', async () => {
|
||||
expect(await localFederatedSourceIds(engine, 'private', 'brain_default')).toEqual(['private', 'default', 'wiki']);
|
||||
});
|
||||
|
||||
test('explicit tiers (--source / env / dotfile) never expand', async () => {
|
||||
expect(await localFederatedSourceIds(engine, 'default', 'flag')).toBeUndefined();
|
||||
expect(await localFederatedSourceIds(engine, 'default', 'env')).toBeUndefined();
|
||||
expect(await localFederatedSourceIds(engine, 'default', 'dotfile')).toBeUndefined();
|
||||
});
|
||||
|
||||
test('single federated source (the resolved one) keeps the scalar fast path', async () => {
|
||||
const solo = { executeRaw: async () => [{ id: 'default' }] } as any;
|
||||
expect(await localFederatedSourceIds(solo, 'default', 'seed_default')).toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
describe('federatedSearchScope — trust + explicitness matrix', () => {
|
||||
test('trusted local + unqualified widens to the federated set', () => {
|
||||
const ctx = ctxOf({ localFederatedSourceIds: ['default', 'wiki'] });
|
||||
expect(federatedSearchScope(ctx)).toEqual({ sourceIds: ['default', 'wiki'] });
|
||||
});
|
||||
|
||||
test('remote caller NEVER widens (fail-closed), even if the field is set', () => {
|
||||
const ctx = ctxOf({ remote: true, localFederatedSourceIds: ['default', 'wiki'] });
|
||||
expect(federatedSearchScope(ctx)).toEqual({ sourceId: 'default' });
|
||||
});
|
||||
|
||||
test('per-call source_id wins over the federated set', () => {
|
||||
const ctx = ctxOf({ localFederatedSourceIds: ['default', 'wiki'] });
|
||||
expect(federatedSearchScope(ctx, 'wiki')).toEqual({ sourceId: 'wiki' });
|
||||
});
|
||||
|
||||
test('per-call __all__ keeps the whole-brain semantics for trusted local', () => {
|
||||
const ctx = ctxOf({ localFederatedSourceIds: ['default', 'wiki'] });
|
||||
expect(federatedSearchScope(ctx, '__all__')).toEqual({});
|
||||
});
|
||||
|
||||
test('a federated OAuth grant wins over the local set', () => {
|
||||
const ctx = ctxOf({
|
||||
localFederatedSourceIds: ['default', 'wiki'],
|
||||
auth: { allowedSources: ['a', 'b'] } as OperationContext['auth'],
|
||||
});
|
||||
expect(federatedSearchScope(ctx)).toEqual({ sourceIds: ['a', 'b'] });
|
||||
});
|
||||
|
||||
test('no local federated set → unchanged scalar scope', () => {
|
||||
expect(federatedSearchScope(ctxOf())).toEqual({ sourceId: 'default' });
|
||||
});
|
||||
});
|
||||
|
||||
describe('search op — unqualified local search spans federated sources', () => {
|
||||
test('federated source results appear; non-federated + archived stay invisible', async () => {
|
||||
const ctx = ctxOf({
|
||||
localFederatedSourceIds: await localFederatedSourceIds(engine, 'default', 'seed_default'),
|
||||
});
|
||||
const results = (await search.handler(ctx, { query: 'zebra telescope' })) as Array<{ slug: string }>;
|
||||
const slugs = results.map((r) => r.slug);
|
||||
expect(slugs).toContain('notes/home');
|
||||
expect(slugs).toContain('wiki/topic'); // pre-#2561 this was missing
|
||||
expect(slugs).not.toContain('private/topic');
|
||||
expect(slugs).not.toContain('old/topic');
|
||||
});
|
||||
|
||||
test('explicit source resolution (no federated set on ctx) stays single-source', async () => {
|
||||
const results = (await search.handler(ctxOf(), { query: 'zebra telescope' })) as Array<{ slug: string }>;
|
||||
const slugs = results.map((r) => r.slug);
|
||||
expect(slugs).toEqual(['notes/home']);
|
||||
});
|
||||
});
|
||||
@@ -16,6 +16,7 @@
|
||||
*/
|
||||
import { describe, test, expect } from 'bun:test';
|
||||
import { createHmac } from 'node:crypto';
|
||||
import { readFileSync } from 'node:fs';
|
||||
import { safeHexEqual } from '../src/core/timing-safe.ts';
|
||||
|
||||
const GITHUB_SECRET = 'super-secret-webhook-key';
|
||||
@@ -123,3 +124,25 @@ describe('Branch ref construction (D5)', () => {
|
||||
expect(pushedRef === `refs/heads/${trackedBranch}`).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe('Webhook sync job extraction contract (#2849)', () => {
|
||||
test('opts into extraction before the pushed commit is consumed', () => {
|
||||
const serveSource = readFileSync(
|
||||
new URL('../src/commands/serve-http.ts', import.meta.url),
|
||||
'utf8',
|
||||
);
|
||||
const routeStart = serveSource.indexOf("'/webhooks/github'");
|
||||
const queueStart = serveSource.indexOf('const job = await queue.add(', routeStart);
|
||||
const responseStart = serveSource.indexOf('res.status(202)', queueStart);
|
||||
expect(routeStart).toBeGreaterThanOrEqual(0);
|
||||
expect(queueStart).toBeGreaterThan(routeStart);
|
||||
expect(responseStart).toBeGreaterThan(queueStart);
|
||||
|
||||
const routeSource = serveSource.slice(queueStart, responseStart);
|
||||
const payload = routeSource.match(
|
||||
/queue\.add\(\s*'sync',\s*\{([\s\S]*?)\}\s*,\s*\{/,
|
||||
);
|
||||
expect(payload).not.toBeNull();
|
||||
expect(payload?.[1]).toMatch(/\bnoExtract:\s*false\b/);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
/**
|
||||
* #2849 — large-sync extract deferral queues an `extract --stale` follow-up.
|
||||
*
|
||||
* performSync's incremental path skips inline link/timeline extraction when
|
||||
* totalChanges > 100 (the #1794 large-sync deferral), leaving
|
||||
* links_extracted_at unstamped. Pre-fix, a standalone sync job (webhook push,
|
||||
* `gbrain sync trigger`) had NOTHING behind it to sweep those pages — the
|
||||
* autopilot cycle's extract phase only walks that cycle's changedSlugs — so a
|
||||
* large webhook push left extraction permanently stale until a manual
|
||||
* `gbrain extract --stale`.
|
||||
*
|
||||
* Pins:
|
||||
* (a) performSync surfaces `extractDeferred: true` on the >100 branch and
|
||||
* leaves the pages unstamped/unlinked.
|
||||
* (b) the `sync` job handler queues an `extract` job with
|
||||
* { stale: true, sourceId? } when extractDeferred is set.
|
||||
* (c) the `extract` handler's stale mode actually sweeps: links created +
|
||||
* watermark stamped (end-to-end recovery, no manual step).
|
||||
*
|
||||
* Marked .serial.test.ts — spawns git subprocesses + shares one PGLite engine.
|
||||
*/
|
||||
|
||||
import { describe, test, expect, beforeAll, afterAll } from 'bun:test';
|
||||
import { mkdtempSync, writeFileSync, rmSync, mkdirSync } from 'fs';
|
||||
import { execSync } from 'child_process';
|
||||
import { tmpdir } from 'os';
|
||||
import { join } from 'path';
|
||||
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
|
||||
import { MinionWorker } from '../src/core/minions/worker.ts';
|
||||
import { MinionQueue } from '../src/core/minions/queue.ts';
|
||||
import { registerBuiltinHandlers } from '../src/commands/jobs.ts';
|
||||
|
||||
let engine: PGLiteEngine;
|
||||
let worker: MinionWorker;
|
||||
let repoPath: string;
|
||||
|
||||
function git(cmd: string): void { execSync(cmd, { cwd: repoPath, stdio: 'pipe' }); }
|
||||
|
||||
describe('#2849 — large sync defers extract and queues a stale sweep', () => {
|
||||
beforeAll(async () => {
|
||||
engine = new PGLiteEngine();
|
||||
await engine.connect({});
|
||||
await engine.initSchema();
|
||||
worker = new MinionWorker(engine, { queue: 'test' });
|
||||
await registerBuiltinHandlers(worker, engine, { quiet: true });
|
||||
|
||||
repoPath = mkdtempSync(join(tmpdir(), 'gbrain-large-defer-'));
|
||||
git('git init');
|
||||
git('git config user.email "t@t.com"');
|
||||
git('git config user.name "T"');
|
||||
mkdirSync(join(repoPath, 'people'), { recursive: true });
|
||||
mkdirSync(join(repoPath, 'notes'), { recursive: true });
|
||||
writeFileSync(join(repoPath, 'people/alice.md'), [
|
||||
'---', 'type: person', 'title: Alice', '---', '', 'Alice is a founder.',
|
||||
].join('\n'));
|
||||
git('git add -A && git commit -m "initial"');
|
||||
|
||||
// Seed: full first sync imports the anchor page + sets last_commit.
|
||||
const { performSync } = await import('../src/commands/sync.ts');
|
||||
await performSync(engine, { repoPath, full: true, noPull: true, noEmbed: true });
|
||||
|
||||
// Second commit: 101 new pages → incremental totalChanges > 100.
|
||||
for (let i = 0; i < 101; i++) {
|
||||
writeFileSync(join(repoPath, `notes/n${i}.md`), [
|
||||
'---', 'type: note', `title: Note ${i}`, '---', '',
|
||||
`[Alice](people/alice) appears in note ${i}.`,
|
||||
].join('\n'));
|
||||
}
|
||||
git('git add -A && git commit -m "add 101 pages"');
|
||||
}, 120_000);
|
||||
|
||||
afterAll(async () => {
|
||||
if (repoPath) rmSync(repoPath, { recursive: true, force: true });
|
||||
if (engine) await engine.disconnect();
|
||||
}, 60_000);
|
||||
|
||||
test('sync handler defers inline extract and queues extract{stale} follow-up; stale sweep recovers', async () => {
|
||||
const syncHandler = (worker as unknown as { handlers: Map<string, (job: unknown) => Promise<unknown>> })
|
||||
.handlers.get('sync');
|
||||
expect(syncHandler).toBeDefined();
|
||||
|
||||
// Same payload shape the webhook submits (minus embed backfill noise).
|
||||
const result = await syncHandler!({
|
||||
data: { repoPath, noExtract: false, noPull: true, auto_embed_backfill: false },
|
||||
signal: { aborted: false },
|
||||
updateProgress: async () => {},
|
||||
}) as { status: string; extractDeferred?: boolean; extract_stale_job_id?: number | null };
|
||||
|
||||
expect(result.status).toBe('synced');
|
||||
// (a) inline extract was deferred, pages left stale.
|
||||
expect(result.extractDeferred).toBe(true);
|
||||
const staleBefore = await engine.countStalePagesForExtraction();
|
||||
expect(staleBefore).toBeGreaterThan(100);
|
||||
expect(await engine.getLinks('notes/n0')).toHaveLength(0);
|
||||
|
||||
// (b) a follow-up extract job with stale:true was queued.
|
||||
expect(result.extract_stale_job_id).toBeGreaterThan(0);
|
||||
const queue = new MinionQueue(engine);
|
||||
const extractJobs = await queue.getJobs({ name: 'extract', limit: 5 });
|
||||
expect(extractJobs.length).toBe(1);
|
||||
expect((extractJobs[0].data as { stale: boolean }).stale).toBe(true);
|
||||
|
||||
// (c) running the extract handler's stale mode recovers: links + stamps.
|
||||
const extractHandler = (worker as unknown as { handlers: Map<string, (job: unknown) => Promise<unknown>> })
|
||||
.handlers.get('extract');
|
||||
await extractHandler!({
|
||||
data: extractJobs[0].data,
|
||||
signal: { aborted: false },
|
||||
updateProgress: async () => {},
|
||||
});
|
||||
const links = await engine.getLinks('notes/n0');
|
||||
expect(links.some(l => l.to_slug === 'people/alice')).toBe(true);
|
||||
const rows = await engine.executeRaw<{ links_extracted_at: string | null }>(
|
||||
`SELECT links_extracted_at FROM pages WHERE slug = 'notes/n0'`,
|
||||
);
|
||||
expect(rows[0]?.links_extracted_at).not.toBeNull();
|
||||
}, 180_000);
|
||||
|
||||
test('sub-threshold sync does NOT set extractDeferred (no spurious follow-up)', async () => {
|
||||
// One more small commit → inline extract path, no deferral.
|
||||
writeFileSync(join(repoPath, 'notes/small.md'), [
|
||||
'---', 'type: note', 'title: Small', '---', '', 'No big deal.',
|
||||
].join('\n'));
|
||||
git('git add -A && git commit -m "one small page"');
|
||||
const { performSync } = await import('../src/commands/sync.ts');
|
||||
const result = await performSync(engine, { repoPath, noPull: true, noEmbed: true });
|
||||
expect(result.status).toBe('synced');
|
||||
expect(result.extractDeferred).toBeFalsy();
|
||||
}, 60_000);
|
||||
});
|
||||
@@ -100,6 +100,10 @@ describe('runSyncTrigger', () => {
|
||||
const job = jobs[0];
|
||||
expect(job.priority).toBe(-10);
|
||||
expect((job.data as { sourceId: string }).sourceId).toBe('default');
|
||||
// #2849: opt in to inline extraction — the standalone sync handler
|
||||
// defaults noExtract to TRUE, which would leave triggered syncs
|
||||
// extraction-stale.
|
||||
expect((job.data as { noExtract: boolean }).noExtract).toBe(false);
|
||||
expect((job.data as { auto_embed_backfill: boolean }).auto_embed_backfill).toBe(true);
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user