mirror of
https://github.com/garrytan/gbrain.git
synced 2026-07-31 04:07:52 +00:00
* wip: federated sync v2 pre-merge snapshot
* v0.40.5.0 Federated Sync v2 — parallel source sync + push triggers + per-source health
Bump VERSION + package.json + CHANGELOG header + migration walkthrough filename
to v0.40.5.0 (claiming the next free slot in the v0.40.x patch series after
master's v0.40.1.0).
What ships (6 components, all behind sync.federated_v2 feature flag default-on):
1. Per-source sync lock — syncLockId(sourceId), phantom-redirect parity
2. Parallel sync --all — pMapAllSettled fan-out, --max-sources N cap
3. embed-backfill minion handler — D2 per-source lock + D6 $10/job budget + D15.1
fire-and-forget submission + D19 source-level cooldown + 24h $25 rolling cap
4. sync trigger CLI + POST /webhooks/github — HMAC-verified (60 req/min/IP),
X-GitHub-Event=push + ref filter against tracked_branch
5. sources status + federation_health doctor — batched GROUP BY pipeline
(4 queries instead of 6×N per-source roundtrips)
6. sources federate/unfederate hook — auto-submit embed-backfill on flip
Correctness fixes (unconditional):
- D21: sync.ts:959 facts backstop now passes sourceId to engine.getPage
- D15.4: redactSourceConfig + CI guard prevent webhook_secret leak
- D15.5: safeHexEqual extracted to src/core/timing-safe.ts
Schema:
- Migration v89 (sources_github_repo_index): partial expression index on
config->>'github_repo' for fast webhook source-lookup
Tests:
- 14 new test files, 112 cases. 4 IRON-RULE regressions pinned (SYNC_LOCK_ID
back-compat, phantom per-source lock, embed-backfill kill+resume,
webhook HMAC prefix-strip). All 9449 unit tests pass.
Caught at test-write time: the webhook handler had a Buffer.from('sha256=...',
'hex') truncation bug — without the prefix-strip, every signature would have
"matched" empty buffers. Pinned by a test/sources-webhook.test.ts IRON-RULE.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
* fix(check-source-config-leak): tighten regex to source-row patterns only
The v0.40.5.0 wave added scripts/check-source-config-leak.sh with a
too-broad pattern (JSON\.stringify\(.*config) that flagged any variable
named 'config' — catching the GLOBAL gbrain config.json serializers in
src/commands/init.ts (status envelopes) and src/core/config.ts (the
config-file write site). On the CI runner without rg installed, the
grep -rE fallback fired correctly and produced 4 false positives that
broke the `verify` script.
Tightened the patterns to specifically match `(source|src|row|s).config`
property access — the actual risk shape (a sources-table row being
serialized whole). The global gbrain config has a different shape and
threat model (file-mode 0o600 at the write site), so it's safe to
exempt at the regex level rather than per-file whitelist.
Also fixed a latent bug: the rg branch used `--include='*.ts'` (grep's
flag, not rg's). rg silently rejected it and CANDIDATES came back empty,
so the local-dev runs (which have rg) would never have caught a real
leak. Now branches on tool availability: `-g '*.ts'` for rg, `--include`
for grep -rE. Both branches verified against a synthetic leak fixture.
Also added init.ts + config.ts to the whitelist as a belt-and-suspenders
since they handle gbrain-global config (not source rows) and could
otherwise reflect-back via regex iteration.
CI: `bun run verify` exit 0 locally with both the original false-positive
fixture (clean repo) and a synthetic leak fixture (correctly caught,
exit 1).
---------
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
231 lines
7.6 KiB
TypeScript
231 lines
7.6 KiB
TypeScript
/**
|
|
* Tests for src/core/embed-stale.ts (v0.40 D15.2).
|
|
*
|
|
* Hermetic — uses an injected `embedFn` so no network call lands. Validates:
|
|
* - empty stale set → done:true, embedded:0
|
|
* - multi-batch run → embed every stale chunk, advance cursor correctly
|
|
* - kill mid-flight (signal.aborted) → aborted:true, partial progress preserved
|
|
* - resume from cursor → picks up where prior call left off (DB predicate)
|
|
* - per-page embedFn throw → logged + skipped, NOT propagated; chunks stay NULL
|
|
*
|
|
* Why PGLite: validates the engine.listStaleChunks/getChunks/upsertChunks
|
|
* roundtrip the helper depends on, not just the loop control flow.
|
|
*/
|
|
import { describe, test, expect, beforeAll, afterAll, beforeEach } from 'bun:test';
|
|
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
|
|
import { resetPgliteState } from './helpers/reset-pglite.ts';
|
|
import { embedStaleForSource } from '../src/core/embed-stale.ts';
|
|
import type { ChunkInput } from '../src/core/types.ts';
|
|
|
|
let engine: PGLiteEngine;
|
|
|
|
beforeAll(async () => {
|
|
engine = new PGLiteEngine();
|
|
await engine.connect({});
|
|
await engine.initSchema();
|
|
}, 30000);
|
|
|
|
afterAll(async () => {
|
|
await engine.disconnect();
|
|
});
|
|
|
|
beforeEach(async () => {
|
|
await resetPgliteState(engine);
|
|
});
|
|
|
|
/** Seed a page with N stale chunks (no embedding) into the default source. */
|
|
async function seedPageWithStaleChunks(slug: string, chunkCount: number): Promise<void> {
|
|
await engine.putPage(slug, {
|
|
type: 'note',
|
|
title: slug,
|
|
compiled_truth: `# ${slug}\n\nseeded`,
|
|
});
|
|
const chunks: ChunkInput[] = Array.from({ length: chunkCount }, (_, i) => ({
|
|
chunk_index: i,
|
|
chunk_text: `chunk ${i} of ${slug}`,
|
|
chunk_source: 'compiled_truth',
|
|
token_count: 4,
|
|
embedding: undefined, // NULL = stale
|
|
}));
|
|
await engine.upsertChunks(slug, chunks);
|
|
}
|
|
|
|
/** Deterministic fake embedder — returns unit-length 1536-dim vectors with
|
|
* first dim = text length, so we can assert specific chunks got embedded. */
|
|
function fakeEmbedFn(texts: string[]): Promise<Float32Array[]> {
|
|
return Promise.resolve(
|
|
texts.map((t) => {
|
|
const v = new Float32Array(1536);
|
|
v[0] = t.length;
|
|
v[1] = 1;
|
|
return v;
|
|
}),
|
|
);
|
|
}
|
|
|
|
describe('embedStaleForSource', () => {
|
|
test('empty stale set returns done:true with zero embedded', async () => {
|
|
const result = await embedStaleForSource(engine, 'default', {
|
|
embedFn: fakeEmbedFn,
|
|
});
|
|
expect(result).toEqual({
|
|
embedded: 0,
|
|
chunksProcessed: 0,
|
|
pagesProcessed: 0,
|
|
lastCursor: null,
|
|
done: true,
|
|
aborted: false,
|
|
});
|
|
});
|
|
|
|
test('embeds every stale chunk across multiple pages in one call', async () => {
|
|
await seedPageWithStaleChunks('a', 5);
|
|
await seedPageWithStaleChunks('b', 3);
|
|
|
|
const result = await embedStaleForSource(engine, 'default', {
|
|
embedFn: fakeEmbedFn,
|
|
});
|
|
expect(result.done).toBe(true);
|
|
expect(result.aborted).toBe(false);
|
|
expect(result.embedded).toBe(8);
|
|
expect(result.pagesProcessed).toBe(2);
|
|
|
|
// Verify DB: zero stale remaining for default.
|
|
const stale = await engine.countStaleChunks({ sourceId: 'default' });
|
|
expect(stale).toBe(0);
|
|
});
|
|
|
|
test('respects batchSize for cursor pagination', async () => {
|
|
await seedPageWithStaleChunks('a', 3);
|
|
await seedPageWithStaleChunks('b', 3);
|
|
let batchCount = 0;
|
|
const result = await embedStaleForSource(engine, 'default', {
|
|
embedFn: fakeEmbedFn,
|
|
batchSize: 2,
|
|
onProgress: () => {
|
|
batchCount++;
|
|
},
|
|
});
|
|
expect(result.embedded).toBe(6);
|
|
// 2-chunk batches across 6 stale rows = at least 3 progress callbacks.
|
|
expect(batchCount).toBeGreaterThanOrEqual(3);
|
|
});
|
|
|
|
test('IRON-RULE: aborted mid-flight → aborted:true, partial progress preserved', async () => {
|
|
await seedPageWithStaleChunks('a', 4);
|
|
await seedPageWithStaleChunks('b', 4);
|
|
await seedPageWithStaleChunks('c', 4);
|
|
const controller = new AbortController();
|
|
// Batch size 4 = one page per batch. concurrency 1 = serialize keys.
|
|
// Abort fires inside embedFn for page 'b', so 'a' lands, 'b' aborts mid-call,
|
|
// and the third batch ('c') never starts.
|
|
const result = await embedStaleForSource(engine, 'default', {
|
|
batchSize: 4,
|
|
concurrency: 1,
|
|
signal: controller.signal,
|
|
embedFn: async (texts) => {
|
|
if (texts.some((t) => t.includes(' of b'))) {
|
|
controller.abort();
|
|
throw new Error('aborted'); // simulates HTTP abort throw
|
|
}
|
|
return fakeEmbedFn(texts);
|
|
},
|
|
});
|
|
expect(result.aborted).toBe(true);
|
|
expect(result.done).toBe(false);
|
|
expect(result.embedded).toBe(4); // only 'a' landed
|
|
// 'b' and 'c' (8 chunks) remain stale
|
|
const stale = await engine.countStaleChunks({ sourceId: 'default' });
|
|
expect(stale).toBe(8);
|
|
});
|
|
|
|
test('IRON-RULE: kill + resume — second call picks up via embedding-IS-NULL predicate', async () => {
|
|
await seedPageWithStaleChunks('a', 4);
|
|
await seedPageWithStaleChunks('b', 4);
|
|
|
|
// First call aborts when 'b' is reached
|
|
const controller = new AbortController();
|
|
const first = await embedStaleForSource(engine, 'default', {
|
|
batchSize: 4,
|
|
concurrency: 1,
|
|
signal: controller.signal,
|
|
embedFn: async (texts) => {
|
|
if (texts.some((t) => t.includes(' of b'))) {
|
|
controller.abort();
|
|
throw new Error('aborted');
|
|
}
|
|
return fakeEmbedFn(texts);
|
|
},
|
|
});
|
|
expect(first.aborted).toBe(true);
|
|
expect(first.embedded).toBe(4); // 'a' landed
|
|
|
|
// Second call with NO cursor — predicate excludes already-embedded chunks
|
|
const second = await embedStaleForSource(engine, 'default', {
|
|
embedFn: fakeEmbedFn,
|
|
});
|
|
expect(second.done).toBe(true);
|
|
expect(first.embedded + second.embedded).toBe(8);
|
|
|
|
const stale = await engine.countStaleChunks({ sourceId: 'default' });
|
|
expect(stale).toBe(0);
|
|
});
|
|
|
|
test('per-page embedFn throw is logged but does NOT propagate', async () => {
|
|
await seedPageWithStaleChunks('good', 2);
|
|
await seedPageWithStaleChunks('bad', 2);
|
|
|
|
let badCount = 0;
|
|
const result = await embedStaleForSource(engine, 'default', {
|
|
embedFn: async (texts) => {
|
|
if (texts.some((t) => t.includes('bad'))) {
|
|
badCount++;
|
|
throw new Error('intentional embed failure');
|
|
}
|
|
return fakeEmbedFn(texts);
|
|
},
|
|
});
|
|
|
|
// The helper itself didn't throw
|
|
expect(result.done).toBe(true);
|
|
expect(badCount).toBe(1);
|
|
|
|
// 'good' chunks got embedded; 'bad' chunks stayed NULL
|
|
expect(result.embedded).toBe(2);
|
|
const stale = await engine.countStaleChunks({ sourceId: 'default' });
|
|
expect(stale).toBe(2);
|
|
});
|
|
|
|
test('source-scoped: does not touch other sources', async () => {
|
|
await engine.executeRaw(
|
|
`INSERT INTO sources (id, name, config) VALUES ('other', 'other', '{"federated":true}'::jsonb) ON CONFLICT (id) DO NOTHING`,
|
|
);
|
|
await seedPageWithStaleChunks('a', 3);
|
|
await engine.putPage('b', {
|
|
type: 'note',
|
|
title: 'b',
|
|
compiled_truth: '# b\n\nseeded',
|
|
}, { sourceId: 'other' });
|
|
await engine.upsertChunks(
|
|
'b',
|
|
Array.from({ length: 3 }, (_, i) => ({
|
|
chunk_index: i,
|
|
chunk_text: `other ${i}`,
|
|
chunk_source: 'compiled_truth',
|
|
token_count: 4,
|
|
embedding: undefined,
|
|
})),
|
|
{ sourceId: 'other' },
|
|
);
|
|
|
|
const result = await embedStaleForSource(engine, 'default', {
|
|
embedFn: fakeEmbedFn,
|
|
});
|
|
expect(result.embedded).toBe(3);
|
|
|
|
// 'other' source still has 3 stale chunks
|
|
const otherStale = await engine.countStaleChunks({ sourceId: 'other' });
|
|
expect(otherStale).toBe(3);
|
|
});
|
|
});
|